【问题标题】:Access record's partition numebr in kafka streams访问记录kafka流中的分区号
【发布时间】: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


【解决方案1】:

请更改

@Override
 public void init(org.apache.kafka.streams.processor.ProcessorContext pc)
 {
                    this.context = context;
}

@Override
 public void init(org.apache.kafka.streams.processor.ProcessorContext pc)
 {
                    this.context = pc;
}

this.context = context // 两者都是一样的,看起来像是错字。

【讨论】:

    猜你喜欢
    • 2021-03-02
    • 1970-01-01
    • 1970-01-01
    • 2020-01-29
    • 1970-01-01
    • 2021-11-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多