【问题标题】:partition count reduced at sink during kafka stream record forward在 kafka 流记录转发期间,接收器的分区计数减少
【发布时间】:2020-01-29 20:46:05
【问题描述】:

我正在使用 kafka 流来处理少量 kafka 记录,我有两个节点,一个用于进行一些转换,另一个是最终接收器。

我的主题是 INTER_TOPIC 和 FINAL_TOPIC 每个有 20 个分区。而我写给 INTER_TOPIC 的生产者正在写键值,而 partition-er 是循环法。

下面是我的转换节点的代码。

public void streamHandler() {

        Properties props = getKafkaProperties();

        StreamsBuilder builder = new StreamsBuilder();

        KStream<String, String> processStream = builder.stream("INTER_TOPIC",
                Consumed.with(Serdes.String(), Serdes.String()));

        //processStream.peek((key,value)->System.out.println("key :"+key+" value :"+value));

        processStream.map((key, value) -> getTransformer().transform(key, value)).filter((key,value)->filteroutFailedRequest(key,value)).to("FINAL_TOPIC", Produced.with(Serdes.String(), Serdes.String()));


        KafkaStreams IStreams = new KafkaStreams(builder.build(), props);

        IStreams.setUncaughtExceptionHandler(new Thread.UncaughtExceptionHandler() {
            @Override
            public void uncaughtException(Thread t, Throw-able e) {

                logger.error("Thread Name :" + t.getName() + " Error while processing:", e);
            }
        });

        IStreams.cleanUp();
        IStreams.start();

        try {
            System.in.read();
        } catch (IOException e) {

            logger.error("Failed streaming ",e);
        }
    }

但我的接收器仅在 2 个分区中获取数据,但我配置了 20 个流线程,并且我验证了我的生产者正在写入所有 20 个分区,如何知道我的转换节点转发到我的 FINAL_TOPIC 的所有 20 个分区

30 Sep 2019 10:39:41,416 INFO  c.j.m.s.StreamHandler [289] [streams-user-61a77203-9afc-4c66-843d-94c20a509793-StreamThread-3] Received
30 Sep 2019 10:39:41,416 INFO  c.j.m.s.StreamHandler [289] [streams-user-61a77203-9afc-4c66-843d-94c20a509793-StreamThread-4] Received
30 Sep 2019 10:39:41,416 INFO  c.j.m.s.StreamHandler [289] [streams-user-61a77203-9afc-4c66-843d-94c20a509793-StreamThread-3] Received
30 Sep 2019 10:39:41,416 INFO  c.j.m.s.StreamHandler [289] [streams-user-61a77203-9afc-4c66-843d-94c20a509793-StreamThread-4] Received
30 Sep 2019 10:40:57,427 INFO  c.j.m.s.StreamHandler [289] [streams-user-61a77203-9afc-4c66-843d-94c20a509793-StreamThread-3] Received
30 Sep 2019 10:40:57,427 INFO  c.j.m.s.StreamHandler [289] [streams-user-61a77203-9afc-4c66-843d-94c20a509793-StreamThread-4] Received
30 Sep 2019 10:40:57,427 INFO  c.j.m.s.StreamHandler [289] [streams-user-61a77203-9afc-4c66-843d-94c20a509793-StreamThread-3] Received
30 Sep 2019 10:40:57,427 INFO  c.j.m.s.StreamHandler [289] [streams-user-61a77203-9afc-4c66-843d-94c20a509793-StreamThread-4] Received

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    partition-er 是循环法

    为什么你认为分区器是循环的?默认情况下,Kafka Streams 根据键应用基于哈希的分区。

    如果要更改默认分区器,可以实现接口StreamPartitioner并通过:

    Produced.with(Serdes.String(), Serdes.String())
            .withStreamPartitioner(...)
    

    【讨论】:

    • 在这种情况下如何启用或添加自定义partitin-er?
    • 更新了我的答案——您可以通过Produced 传递自定义分区器。
    • props.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, CustomPartitioner.class); ……这也行吗?这段代码加了一点
    • 你能举一些例子来说明如何写 withStreamPartitioner(...) 吗?
    • 也许默认实现会有所帮助(不知道其他示例):github.com/apache/kafka/blob/trunk/streams/src/main/java/org/…
    猜你喜欢
    • 1970-01-01
    • 2022-01-22
    • 2018-01-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-05-28
    • 1970-01-01
    相关资源
    最近更新 更多