【问题标题】:Kafka Connect - Transformes rename field only if it existKafka Connect - 仅在存在时才转换重命名字段
【发布时间】:2020-11-26 23:55:52
【问题描述】:

我有一个用于多个主题(topic_a、topic_b、topic_c)的 S3 接收器连接器,并且 topic_a 有字段 created_date 和 topic_b、topic_c 有 creation_date。我使用下面的transforms.RenameField.renames 重命名了该字段 (created_date:creation_date) 但由于唯一的 topic_a 有 created_date 而其他没有,连接器失败了。

我想将所有消息(来自具有单个连接器的所有主题)移动到带有 creation_date 的 s3 中(如果存在,则将 created_date 重命名为 creation_date),但我无法找出正则表达式或转换器来重命名该字段(如果它存在)针对特定主题。

   "config":{
      "connector.class":"io.confluent.connect.s3.S3SinkConnector",
      "errors.log.include.messages":"true",
      "s3.region":"eu-west-1",
      "topics.dir":"dir",
      "flush.size":"5",
      "tasks.max":"2",
      "s3.part.size":"5242880",
      "timezone":"UTC",
      "locale":"en",
      "format.class":"io.confluent.connect.s3.format.json.JsonFormat",
      "errors.log.enable":"true",
      "s3.bucket.name":"bucket",
      "topics": "topic_a, topic_b, topic_c",
      "s3.compression.type":"gzip",
      "partitioner.class":"io.confluent.connect.storage.partitioner.DailyPartitioner",
      "name":"NAME",
      "storage.class":"io.confluent.connect.s3.storage.S3Storage",
      "key.converter.schemas.enable":"true",
      "key.converter":"org.apache.kafka.connect.storage.StringConverter",
      "value.converter.schemas.enable":"true",
      "value.converter":"io.confluent.connect.avro.AvroConverter",
      "value.converter.schema.registry.url":"https://schemaregistry.com",
      "enhanced.avro.schema.support": "true",
      "transforms": "RenameField",
      "transforms.RenameField.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
      "transforms.RenameField.renames": "created_date:creation_date"
   }

【问题讨论】:

    标签: apache-kafka apache-kafka-connect


    【解决方案1】:

    只有 topic_a 有 created_date 而其他人没有,

    然后您将使用单独的连接器。一个带有变换,所有主题都带有字段,然后另一个没有变换。

    来自所有具有单个连接器的主题

    这不能很好地扩展。您正在制作有限的消费者线程和一个消费者组来一次阅读多个主题。多个连接器会更好地分配负载。

    【讨论】:

    • 谢谢@onecricketeer,我知道最好的解决方案是创建单独的连接器,但分配的负载不多,这就是我寻找一个连接器的原因。如果该字段存在,您是否有任何想法重命名该字段,否则按原样复制消息。
    • 可惜connect本身不能轻易做这样的逻辑。您必须编写 KStreams 才能在该级别上做某事
    • 谢谢,我已经为方便主题使用了单独的连接器。
    猜你喜欢
    • 2020-11-20
    • 2019-08-22
    • 2017-06-08
    • 1970-01-01
    • 2018-06-14
    • 1970-01-01
    • 1970-01-01
    • 2019-08-09
    • 2018-04-19
    相关资源
    最近更新 更多