【问题标题】:Kafka replicator: ConsumerTimestampsInterceptor with kafka Streams?Kafka复制器:带有kafka Streams的ConsumerTimestampsInterceptor?
【发布时间】:2021-06-11 19:22:13
【问题描述】:

我们正在尝试在 2 个数据中心之间复制我们的偏移量。对于单个消费者来说真的很容易,只需添加:

consumer.interceptor.classes=io.confluent.connect.replicator.offsets.ConsumerTimestampsInterceptor

现在我们有一个使用 kafka-streams 的应用程序。在绑定多个事物之后,我们无法像之前那样复制偏移量。 例如我们也尝试过:

kafka.streams.properties.consumer.interceptor.classes=io.confluent.connect.replicator.offsets.ConsumerTimestampsInterceptor

但没有运气 谢谢!

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    试着用代码而不是属性文件来做

    例如生产者

    config.put(
        StreamsConfig.PRODUCER_PREFIX + ProducerConfig.INTERCEPTOR_CLASSES_CONFIG, 
        ConsumerTimestampsInterceptor.class.getName()
    );
    

    或者,确保 kafka.streams.properties 是创建所有 StreamsConfig 属性的正确属性前缀

    【讨论】:

    • 我有:ClassCastException: class io.confluent.connect.replicator.offsets.ConsumerTimestampsInterceptor ?
    • 我的回答没有做任何转换,所以你能用完整的堆栈跟踪编辑你的问题吗?
    【解决方案2】:

    最后 Kafka Streams 删除了 group.id 参数。但是复制器仍然需要它。

        override fun configure(configs: Map<String, *>?) {
            val newConfigs = configs!! + mapOf("group.id" to configs?.get("customgroupid"))
            super.configure(newConfigs)
        }
    }
    

    然后将其添加到您的 .properties 文件中:

    kafka.streams.consumer.customgroupid=${kafka.streams.group.id} # customgroupid is not removed
    kafka.streams.consumer.interceptor.classes=com.myjob.mymodule.ConsumerTimestampsInterceptorConfigurator
    

    【讨论】:

      猜你喜欢
      • 2017-03-15
      • 2018-07-15
      • 1970-01-01
      • 2018-01-19
      • 2019-01-31
      • 2020-04-25
      • 2017-02-21
      • 2019-09-07
      • 2019-03-30
      相关资源
      最近更新 更多