【问题标题】:Kafka connection transformations. Parse input string and get a record keyKafka 连接转换。解析输入字符串并获取记录键
【发布时间】:2018-07-16 16:20:18
【问题描述】:

我使用一个简单的文件源阅读器

connector.class=org.apache.kafka.connect.file.FileStreamSourceConnector
tasks.max=1

文件内容是每一行中的一个简单 JSON 对象。我发现有一种方法可以替换记录键并使用转换来执行此操作,例如

# Add the `id` field as the key using Simple Message Transformations
transforms=InsertKey

# `ValueToKey`: push an object of one of the column fields (`id`) into the key
transforms.InsertKey.type=org.apache.kafka.connect.transforms.ValueToKey
transforms.InsertKey.fields=ip

但是我遇到了一个错误

仅支持 [将字段从值复制到键] 的 Struct 对象, 找到:java.lang.String

有没有办法像使用 Flume 和 regex_extractor 一样解析字符串 json 并从那里获取密钥?

【问题讨论】:

  • 你的键是一个字符串,而不是一个结构。您如何期望将某些内容“插入”到字符串中?
  • 我已将配置替换为 transforms.ReplaceKey.type=org.apache.kafka.connect.transforms.ReplaceField$Key transforms.ReplaceKey.whitelist=ip 但它仍然会产生错误Only Map objects supported in absence of schema
  • 我发现如果源没有模式是不可能的。它必须是另一个支持生成结构化模型的插件。
  • @SergeiGrigorev,在没有源代码的情况下,您是否在以后的任何时候发现了这个问题?

标签: json apache-kafka apache-kafka-connect


【解决方案1】:

有一种方法可以替换记录键

有一个单独的转换称为org.apache.kafka.connect.transforms.ReplaceField$Key

InsertKey 将获取一个值并尝试插入到 Struct/Map 中,但您似乎使用的是字符串键

【讨论】:

  • 我改成了 ReplaceField$Key,但我仍然有一个问题,我的输入值没有架构或类似 org.apache.kafka.connect.errors.DataException: Only Map objects在没有 [字段替换] 架构的情况下支持,发现:null
  • org.apache.kafka.connect.errors.DataException: Only Map objects supported in absence of schema for [field replacement], found: null at org.apache.kafka.connect.transforms.util.Requirements.requireMap(Requirements.java:38) at org.apache.kafka.connect.transforms.ReplaceField.applySchemaless(ReplaceField.java:133) at org.apache.kafka.connect.transforms.ReplaceField.apply(ReplaceField.java:126) at org.apache.kafka.connect.runtime.TransformationChain.apply(TransformationChain.java:38)
  • 嗯。我觉得没有模式,它只是不知道如何解析你的对象。这就是为什么 Avro 比简单的 JSON 更受欢迎的原因
  • 但我有输入日志文件,其中每一行都是一个 json 字符串。在水槽中,我可以在拦截器中使用 regex_extractor。我尝试对 kafka connect 做同样的事情。所以......我会尝试使用其他支持 json 的插件来做到这一点。无论如何,谢谢你的帮助
  • Confluent 文档明确指出 Filestream 源仅用于开始连接,而不是生产用例。 Flume、Fluentd、Filebeat 都是消费日志文件到 Kafka 的绝佳选择
【解决方案2】:

在 SourceConnector 上使用转换时,转换是在 SourceConnector.poll() 返回的 List<SourceRecord> 上完成的

在您的情况下,FileStreamSourceConnector 读取文件的行并将每一行作为字符串放入 SourceRecord 对象中。因此,当变换得到SourceRecord时,它只把它看作一个String,并不知道对象的结构。

为了解决这个问题,

  1. 要么修改FileStreamSourceConnector 代码,以便它返回SourceRecord 和输入json 字符串的有效StructSchema。您可以为此使用 Kafka 的 SchemaBuilder 类。
  2. 或者,如果您在接收器连接器中使用此数据,您可以通过在接收器连接器上设置以下配置,让 KafkaConnect 将其转换为 JSON,然后在接收器连接器上执行转换。

"value.converter":"org.apache.kafka.connect.json.JsonConverter"
"value.converter.schemas.enable": "false"

如果您选择第二个选项,请不要忘记将这些配置放在您的 SourceConnector 上。

"value.convertor":"org.apache.kafka.connect.storage.StringConverter"
"value.converter.schemas.enable": "false"

【讨论】:

    猜你喜欢
    • 2021-10-16
    • 1970-01-01
    • 1970-01-01
    • 2011-11-01
    • 2012-01-21
    • 1970-01-01
    • 2015-03-09
    • 1970-01-01
    • 2020-08-20
    相关资源
    最近更新 更多