【问题标题】:Creating kafka stream API for JSON's为 JSON 创建 kafka 流 API
【发布时间】:2017-11-27 09:15:02
【问题描述】:

我正在尝试编写一个 kafka 流代码,用于将 JSON 数组转换为 JSON 元素......因为我是 kafka 流的新手,任何人都可以帮助我编写代码......就像 kstream 和 ktable 中应该有的一样。 . 我的输入流将采用以下格式

[
 {"timestamp":"2017-10-24T12:44:09.359126933+05:30","data":0,"unit":""},
 {"timestamp":"2017-10-24T12:44:09.359175426+05:30","data":1,"unit":""}
]

[
 {"timestamp":"2017-10-24T12:44:09.359126933+05:30","data":2,"unit":""},
 {"timestamp":"2017-10-24T12:44:09.359175426+05:30","data":3,"unit":""}
]

我的输出必须是格式

{"timestamp":"2017-10-24T12:44:09.359126933+05:30","data":0,"unit":""}
{"timestamp":"2017-10-24T12:44:09.359175426+05:30","data":1,"unit":""}
{"timestamp":"2017-10-24T12:44:09.359126933+05:30","data":2,"unit":""}
{"timestamp":"2017-10-24T12:44:09.359175426+05:30","data":3,"unit":""}

谁能帮我写代码??

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    如果你想使用 Kafka Streams,你可以使用flatMap()。类似的东西

    // using new 1.0 API
    StreamsBuilder builder = new StreamsBuilder();
    builer.stream("topic").flatMap(...).to("output-topic");
    

    查看示例和文档了解更多详情:

    【讨论】:

    • 所以如果我使用平面地图,我可以将 Json 数组转换为 Json 对象,就像我在上面的问题中显示的那样?
    • 我在 KStreams 中的 应该是什么?就像我需要为此指定什么数据类型..?非常感谢你帮我做这件事..我被困在这一点上......
    • 是的。您在 flatMap 中进行转换。对于您的输入流, 类型必须与您在输入主题中使用的类型相匹配。对于输出类型,你可以选择你想写的任何类型作为结果类型。
    • 我的输入将采用 Json Araays 的形式,那么我的键值对可能是什么?我的输出应该是 Json 对象......在平面图中我应该使用 split 还是我可以使用什么?
    • 如果你不使用密钥,它可以是null,你可以只使用byte[]一个密钥类型。如果您不知道如何拆分 JSON 数组,请阅读 JSON 教程...
    【解决方案2】:

    在 Python 中...

    from kafka import KafkaConsumer
    consumer = KafkaConsumer('topicName')
    for message in consumer:
      print(message)
    

    在 KafkaConsumer 中指定 bootstrap_servers 参数。

    对于Java 看cloudkarafka,真的不错:

    KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
    consumer.subscribe(Arrays.asList(topic));
    while (true) {
        ConsumerRecords<String, String> records = consumer.poll(100);
        for (ConsumerRecord<String, String> record : records)
            System.out.printf("msg = %s\n", record.value());
        }
    }
    

    【讨论】:

    • 我需要在 java..如果你知道你能帮我吗?
    • 完成了,在我看来,Scala 更好;)
    • 问题是关于 Kafka Streams API——而不是普通的 Kafka Consumer。
    • @MatthiasJ.Sax :如何将此字符串转换为 json 值并发送到另一个主题?你有例子吗?
    猜你喜欢
    • 2018-08-17
    • 1970-01-01
    • 2018-12-21
    • 2018-12-13
    • 1970-01-01
    • 2019-10-01
    • 2018-12-24
    • 1970-01-01
    • 2019-11-16
    相关资源
    最近更新 更多