Jas*_*son 1 apache-kafka apache-kafka-streams
我们已经开始试验 Kafka,看看它是否可以用来聚合我们的应用程序数据。我认为我们的用例与 Kafka 流匹配,但我们不确定我们是否正确使用该工具。我们构建的概念验证似乎按设计工作,我不确定我们是否正确使用了 API。
我们的概念证明是使用 kafka 流来保存有关输出主题中程序的信息的运行记录,例如
{
"numberActive": 0,
"numberInactive": 0,
"lastLogin": "01-01-1970T00:00:00Z"
}
Run Code Online (Sandbox Code Playgroud)
计算计数很容易,它本质上是根据输入主题和输出字段执行比较和交换(CAS)操作。
本地状态包含给定密钥的最新程序。我们针对状态存储加入输入流,并使用 TransformSupplier 运行 CAS 操作,该操作使用 TransformSupplier 将数据显式写入状态存储
context.put(...)
context.commit();
Run Code Online (Sandbox Code Playgroud)
这是对当地国营商店的适当使用吗?是否还有另一种方法可以在主题中保持状态运行记录?
小智 5
您的设计对我来说听起来很正确(我假设您使用的是 PAPI 而不是 Streams DSL),您正在一个流中读取,在状态存储与运算符关联的流上调用 Transform() 。由于您的更新逻辑似乎仅依赖于键,因此可以通过基于键分区的 Streams 库实现令人尴尬的并行化。
需要注意的一件事是,您似乎在每次 put 调用后都调用“context.commit()”,这不是推荐的模式。这是因为commit()操作是一个相当繁重的调用,涉及刷新状态存储、向 Kafka 代理发送提交偏移请求等,在每个调用上调用它都会导致非常低的吞吐量。建议仅在处理完一堆记录后才调用 commit(),或者您可以仅依靠 Streams 配置“commit.interval.ms”来依赖 Streams 库仅在每个时间间隔后在内部调用 commit() 。请注意,这不会影响正常关闭时的处理语义,因为关闭时 Streams 将始终强制执行 commit() 调用。
| 归档时间: |
|
| 查看次数: |
1433 次 |
| 最近记录: |