【问题标题】:Kafka Stream - Filter by client_idKafka Stream - 按 client_id 过滤
【发布时间】:2020-12-16 18:34:01
【问题描述】:

我正在使用 Kafka Stream 创建一个仅包含特定于 client_id 的数据的 ktable,这不是主题键。我是 Kafka Streams 的新手,这看起来很简单,但我对社区中可用的多个示例感到有些困惑,这些示例非常好。

我正在尝试获取具有 client_id=0123456 的 inputTopic 数据。在下面的 KSQL 中将类似于命令:

CREATE STREAM TOPIC1_CLIENT1 AS
SELECT * FROM TOPIC1
WHERE client_id= '0123456'
EMIT CHANGES;

下面我试图重现相同的行为。有人可以告诉我在下面做错了什么吗?它没有像我预期的那样过滤。

        final KStream<String, String> stream = builder.stream(inputTopic, Consumed.with(stringSerde, stringSerde));
        final KTable<String, String> convertedTable = stream.filter((client_id,v) -> v.equals("0123456")).toTable(Materialized.as("stream-converted-to-table"));
        stream.to(streamsOutputTopic, Produced.with(stringSerde, stringSerde));
        convertedTable.toStream().to(tableOutputTopic, Produced.with(stringSerde, stringSerde));

【问题讨论】:

    标签: apache-kafka stream ksqldb ktable


    【解决方案1】:

    v 是消息的整个值。要在 KSQL 中命名字段,流上有一个关联的模式,例如数据是 JSON 还是 Avro,这意味着 clientid 只是值的一部分

    【讨论】:

    • 所以我只是想在我的 Stream 中创建和关联架构?感谢您查看此内容。
    • 如果可以,请您参考一些文档。关于它?
    • Kafka Streams 没有模式。问题是您使用的是字符串 serde。我不知道你的源数据格式是什么,但假设你有一个像{"client_id": 123} 这样的记录,那么整个字符串当然不会等于 123...你可以参考流上的 Apache 或 Confluent 网站
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-01-19
    • 2021-03-16
    • 2018-10-20
    • 1970-01-01
    • 2018-12-29
    相关资源
    最近更新 更多