【问题标题】:How to read from a string delimited kafka topic then convert it to json?如何从字符串分隔的kafka主题中读取然后将其转换为json?
【发布时间】:2020-04-03 07:50:25
【问题描述】:

我目前正在开发一个 POC 项目,该项目使用 Kafka Connect 从 Kafka 主题中读取数据以写入数据库。我知道 JDBC 接收器连接器需要架构才能工作。然而,我们所有的 kafka 主题都是字符串分隔的。除了使用 json 或 avro 创建一个新主题之外,我还计划创建一个将字符串转换为 json 的 API,我可以尝试任何可能的解决方案吗?

【问题讨论】:

  • 我建议使用流处理器,而不是“API”

标签: apache-kafka apache-kafka-connect


【解决方案1】:

您的问题的解决方案是创建一个自定义转换器。不幸的是,没有博客文章或在线资源说明如何逐步完成此操作,但您可以查看StringConverter 是如何构建的。这是一个非常简单的 API,它使用序列化器和反序列化器来转换数据。

L.E:

转换器使用序列化器和反序列化器,因此您唯一应该做的就是创建自己的:

public class StringConverter implements Converter {

   private final CustomStringSerializer serializer = new CustomStringSerializer();
   private final CustomStringSerializer deserializer = new CustomStringSerializer();

   --- boilerplate code resulted by implementing the converter Interface ---

      @Override
    public byte[] fromConnectData(String topic, Schema schema, Object value) {
        try {
            return serializer.serialize(topic, value == null ? null : value.toString());
        } catch (SerializationException e) {
            throw new DataException("Failed to serialize to a string: ", e);
        }
    }

    @Override
    public SchemaAndValue toConnectData(String topic, byte[] value) {
        try {
            return new SchemaAndValue(createYourOwnSchema(topic, value), deserializer.deserialize(topic, value));
        } catch (SerializationException e) {
            throw new DataException("Failed to deserialize string: ", e);
        }
    }

}

您会注意到toConnectData 正在返回一个SchemaAndValue 对象,实际上它正是您需要自定义的。您可以创建自己的 createYourOwnSchema 方法,该方法将根据主题(和值?)创建架构,然后您拥有反序列化器来反序列化您的数据。

【讨论】:

  • KSQL 也可以,假设有一个字符串拆分函数
  • 非常感谢您的意见。如果这看起来很明显,我很抱歉,但我只需要保证可以创建一个使用字符串分隔的 kafka 主题的自定义转换器,然后我会添加一个方法来将模式包含到每条消息中。我做对了吗?感谢您的所有帮助!
  • 简短回答:是的!我添加了一些关于如何使用转换器来实现的细节。
  • 帮助很大!非常感谢你做的这些。如果我错了,请纠正我。如果我做对了,转换器是整个 kafka 连接框架的一部分,对吗?如果我要创建一个自定义转换器,我是否需要创建一个从消费到 kafka 主题到转换器的使用的整个框架?对不起,请耐心等待。
  • @Isabelle 无需道歉! 1. 是的,转换器是 Kafka Connect 框架的一部分。 2. 不,因为转换器是 Kafka Connect 框架的一部分,你唯一需要做的就是将它安装在你的工作人员上(就像你安装一个新的连接器,链接在这里:docs.confluent.io/current/connect/managing/…)然后,在部署期间您的连接器,您可以指定它应该使用哪个转换器:例如 ... "value.converter": "com.custom.converter.MyCustomConverter" ....
【解决方案2】:

如果我了解您的用例,我肯定会尝试使用 Kafka Streams 来处理您的分隔字符串主题并将其放入 JSON 主题中。然后,您可以将此 JSON 主题用作数据库的输入。

使用 Kafka Streams,您可以映射主题的每个字符串记录并将逻辑应用到其中以生成包含数据的 JSON。处理后的记录可以像 JSON 格式的字符串或什至使用 Kafka JSON Serdes (Kafka Streams DatatTypes) 的 JSON 类型一样沉入您的结果主题。

这种方法还增加了输入数据的灵活性,您可以处理、转换或忽略字符串的每个分隔字段。

如果您想开始使用,请查看Kafka DemoConfluent Kafka-Streams examples,了解更复杂的用例。

【讨论】:

  • 输入好!我会调查一下。非常感谢您接受我的查询。但是,您知道我有什么方法可以通过不创建/摄取到另一个主题而是直接将其插入数据库来做到这一点?
  • 不幸的是,我认为不可能将数据库用作 Kafka Streams 应用程序的输出。 Kafka Streams 旨在使用 Source 处理器处理存储在 Kafka Topic 中的数据,并通过 Sink 处理器使用您的 Kafka 集群作为输出。如果您想将数据从 Kafka 移动到数据库,您应该使用 Kafka Connect,但我建议您先使用 Kafka Streams 来容纳您的数据。此示例总结了工作流程:confluent.io/blog/hello-world-kafka-connect-kafka-streams
【解决方案3】:

让我们创建 Kafka Streams 应用程序。您将使用数据作为 KStream<String, String>(假设密钥也是String) 而不是在 Kafka Stream 应用程序中将其转换为 KStream<String,YourAvroClass> 并生成消息以针对 Avro 主题。

KStream<String,YourAvroClass> avroKStream ....
avroKStream.to("avro-output-topic", Produced.with(Serdes.String(), yourAvroClassSerde));

【讨论】:

    猜你喜欢
    • 2021-08-21
    • 2018-05-07
    • 1970-01-01
    • 2016-10-01
    • 1970-01-01
    • 1970-01-01
    • 2019-11-18
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多