【发布时间】:2020-03-10 12:07:32
【问题描述】:
我已经给出了从主题、处理器和接收器到其他主题的拓扑
StoreBuilder storeBuilder = Stores.keyValueStoreBuilder(
Stores.persistentKeyValueStore("store"),
Serdes.String(),
Serdes.String());
Topology topology = new Topology();
topology.addSource("incoming", Serdes.String().deserializer(), Serdes.String().deserializer(), "topic");
topology.addProcessor("incoming_first", () -> new MyProcessor(), "incoming");
topology.addStateStore(storeBuilder, "incoming_first");
topology.addSink("sink", "sink", "incoming_first"),
public class MyProcessor implements Processor<String, String> {
private ProcessorContext context;
private KeyValueStore<String, String> stateStore;
@Override
public void init(ProcessorContext context) {
this.context = context;
this.stateStore = (KeyValueStore<String, String>) context.getStateStore("store");
}
@Override
public void process(String key, String value) {
stateStore.put(key, value);
....
throw new RuntimeException();
....
context.forward(); //forward to sink
}
@Override
public void close() {
}
}
我的问题是如何处理写入状态存储后处理器发生异常的情况。 Kafka是否有一些带有状态存储回滚的错误处理机制来重新处理消息或将其转发到错误主题?
目前,在没有任何处理的情况下,我的应用程序完全死机,我需要重新启动它。 此外,如果我添加一些 try-catch 消息,则标识为 ok 并且我的状态存储被更新,并且消息被发送到 changelog 主题。
我需要一些状态存储的回滚机制吗?
https://issues.apache.org/jira/browse/KAFKA-7192 KIP 说如果发生异常,状态存储不应该使用 EOS 处理,但这仅适用于我的整个应用程序死亡的情况。
提前致谢!
【问题讨论】:
标签: apache-kafka apache-kafka-streams