【发布时间】:2022-01-29 06:45:09
【问题描述】:
我正在使用带有 spring cloud stream kafka-streams binder 的 Kafka 2.6。我想在我的 Kafka 流应用程序中访问记录头、分区号等。我阅读了有关使用处理器 API、使用 ProcessorContext 等的信息。但是每次 ProcessorContext 对象都为空。
下面是代码
@StreamListener(Bindings.input)
@SendTo(Bindings.output)
public KStream<String, String> process(KStream<String, String> input)
{
return input.transform(new TransformerSupplier<String, String, KeyValue<String, String>>()
{
public Transformer<String, String, KeyValue<String, String>> get()
{
return new Transformer<String, String, KeyValue<String, String>>()
{
private int total = 0;
ProcessorContext context;
@Override
public void close() {
}
@Override
public void init(org.apache.kafka.streams.processor.ProcessorContext pc)
{
this.context = context;
}
@Override
public KeyValue<String, String> transform(String k, String v)
{
System.out.println("ProcessorContext: "+this.context);
System.out.println("value: "+v);
return new KeyValue<>(k, v);
}
};
}
});
}
在此代码中,ProcessorContext 始终打印为 null。我还尝试使用 ListenerContainerCustomizer 进行弹簧启动。但这也行不通
@Bean
ListenerContainerCustomizer<AbstractMessageListenerContainer<?, ?>> customizer()
{
return (container, dest, group) ->
{
container.setRecordInterceptor(record ->
{
System.out.println(">>>> Received record, checking headers");
Headers headers = record.headers();
System.out.println(">>>> Header length: " + headers.toArray().length);
for (Header header : headers)
{
if (header.key().equalsIgnoreCase("eventtype"))
{
String value = String.valueOf(header.value());
if (!value.equalsIgnoreCase("PUBLISHED"))
{
System.out.println("Event type from header not PUBLISHED, skipping record");
return null;
}
}
}
System.out.println("Processing record");
return record;
});
};
}
我打印了在上面看到的 bean 注册的 bean 列表。但它永远不会奏效。 无论如何,我需要第一种方法来工作,因为我喜欢使用分区号运行一些业务逻辑。
请帮助因为很多天而严重卡住。
【问题讨论】:
-
这里有些可疑 - @Override public void init(org.apache.kafka.streams.processor.ProcessorContext pc) { this.context = context; } 应该 this.context = context;成为this.context = pc;而不是上下文
-
天哪!我多么愚蠢。非常感谢伙计。请将其发布为答案,以便我接受
-
最好完全删除您的问题。答案不会解决任何技术任务。只是一个错字并不能证明一个完整的答案。
标签: java apache-kafka apache-kafka-streams spring-cloud-stream kafka-streams-binder