使用Confluent Replicator时,如何实现支持Headers的Kafka Converter?
在Confluent Replicator中实现支持Headers的Kafka Converter问题
我需要实现一个依赖Kafka消息Headers的自定义Kafka Connect Converter,根据Kafka Connect的Converter接口定义,带Headers参数的toConnectData方法应为系统优先调用的版本,不带参数的是为向后兼容保留的旧方法:
byte[] fromConnectData(String topic, Schema schema, Object value); byte[] fromConnectData(String topic, Headers headers, Schema schema, Object value)
但实际使用Confluent Replicator(镜像版本confluentinc/cp-enterprise-replicator:7.7.0)运行时,系统却调用了不带Headers的旧方法,导致业务功能无法完成。
我的Converter实现示例如下:
package com.example; import org.apache.kafka.common.header.Headers; import org.apache.kafka.connect.data.Schema; import org.apache.kafka.connect.data.SchemaAndValue; import org.apache.kafka.connect.storage.Converter; public class ExampleConverter implements Converter { ... @Override public SchemaAndValue toConnectData(String topic, byte[] value) { throw new RuntimeException("headers not supplied, these are required in order to decrypt"); } @Override public SchemaAndValue toConnectData(String topic, Headers headers, byte[] value) { return new SchemaAndValue(Schema.BYTES_SCHEMA, null); } }
运行时出现以下错误,确认调用了不带Headers的方法:
java.lang.RuntimeException: headers not supplied, these are required in order to decrypt at com.example.ExampleConverter.toConnectData(ExampleConverter.java:50) at io.confluent.connect.replicator.ReplicatorSourceTask.convertKeyValue(ReplicatorSourceTask.java:637) at io.confluent.connect.replicator.ReplicatorSourceTask.poll(ReplicatorSourceTask.java:536) at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.poll(AbstractWorkerSourceTask.java:488) at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.execute(AbstractWorkerSourceTask.java:360) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:229) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:284) at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.run(AbstractWorkerSourceTask.java:80) at org.apache.kafka.connect.runtime.isolation.Plugins.lambda$withClassLoader$7(Plugins.java:339) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) at java.base/java.lang.Thread.run(Thread.java:840)
请问哪里操作有误?求相关建议。
内容的提问来源于stack exchange,提问作者Peter McIntyre
相关产品推荐
相关产品推荐

