【发布时间】:2018-02-02 06:45:16
【问题描述】:
如果我的 kafka 流应用程序中有一个共享变量,并且在处理代码中由多个线程更新,它是如何处理的?我是否必须使该共享变量线程安全,或者 Kafka 流库如何处理?在文档的某处,我读到运行 Kafka 流应用程序时不需要在线程之间进行协调。例如,这是一个伪代码:
KStream<byte[], byte[]> input = ...;
int counter = 0;
KStream<byte[], byte[]>[] processed = input.map(
(k, v) -> {
....
....
//update counter by multiple threads.
);
如果此代码由来自同一应用实例的多个流任务执行,计数器会发生什么情况?变量“已处理”如何,因为它也可以由多个线程更新?这需要在普通 Java 场景中进行某种同步。我很好奇这是否由 Kafka 流库处理。
谢谢!
【问题讨论】: