【问题标题】:Drop a message Kafka Streams Topology发送消息 Kafka Streams Topology
【发布时间】:2021-11-30 12:01:13
【问题描述】:

我想知道是否有办法从流拓扑中删除记录/消息?

我有如下设置:

 builder.stream("my-source-topic")
                .map(CustomMapper)
                .mapValues(CustomValueMapper)
                .filterNot(CustomFilter)
                .transformValues(CustomValueTransformer)
                .toStream()

每个 CustomMapper/CustomFilter 等都会覆盖它们各自的应用/转换方法,它们可能如下所示,如上所述,错误可能无法恢复,这是一个好的解决方案,这些消息将被手动处理并写入相应的日志。 假设在第一个映射期间发生不可恢复的错误,我现在如何防止后续阶段甚至处理记录,我想停止处理此记录并移至下一条记录。

@Override
    public V transform(K readOnlyKey, V value) {
        try {
        // do some logic
        } catch(Exception e){
            // process error - this might be unrecoverable.
            
            dropRecord(); // this is what i would be looking for if possible
        }
    }

我可以杀死线程并让 customUncaughtExceptionHandler 重新调度不会提交偏移量的线程,因此尝试再次处理错误记录。

为传递的对象创建包装器需要在每个处理步骤中添加检查以查看记录是否仍然有效。

在每个处理步骤之前添加 .branch() 也需要大量返工。

【问题讨论】:

  • 您可以让地图在其输出中添加一个“无效”标志,并在地图之后有一个简单的过滤器来过滤掉这些对象。
  • 是的,但是我必须在每个处理步骤之后添加过滤器,我还必须创建一个包含要标记为无效字段的包装类。

标签: java apache-kafka-streams


【解决方案1】:

您只需返回null 即可在Transformer 中发送消息。请参阅Transformer#transform 的 Javadoc。 所以你的例子是:

    @Override
    public V transform(K readOnlyKey, V value) {
        try {
        // do some logic
        } catch(Exception e){
            // process error - this might be unrecoverable.
            
            return null;
        }
    }

请注意,您目前只能在 Transformer 中执行此操作,但不能在 ValueTransformer 中执行此操作。

【讨论】:

  • 欢迎来到 Stack Overflow,感谢您发布有用的答案!
猜你喜欢
  • 2019-02-25
  • 1970-01-01
  • 1970-01-01
  • 2021-07-13
  • 2020-03-24
  • 2019-10-11
  • 1970-01-01
  • 2020-09-18
  • 2017-02-23
相关资源
最近更新 更多