【问题标题】:kafka message key as key field/column in HDFSkafka 消息键作为 HDFS 中的关键字段/列
【发布时间】:2021-04-12 04:13:56
【问题描述】:

所以我在我的 MySQL 源连接器中使用 debezium key.field.name 将一个字段添加到我的主题中。

登陆主题后消息如下所示。

{"id":20,"__PKtableowner":"reviewDB.review.search_user_02"}:{"before":null,"after":{"search_user_all_shards.Value":{"id":20,"name":{"string":"oliver"},"email":{"string":"bob@abc.com"},"department":{"string":"sales"},"modified":"2021-01-06T09:27:02Z"}},"source":{"version":"1.0.2.Final","connector":"mysql","name":"reviewDB","ts_ms":1609925222000,"snapshot":{"string":"false"},"db":"review","table":{"string":"search_user_02"},"server_id":1,"gtid":null,"file":"binlog.000002","pos":14630,"row":0,"thread":{"long":13},"query":null},"op":"c","ts_ms":{"long":1609925222956}}

在哪里, 关键是

{"id":20,"__PKtableowner":"reviewDB.review.search_user_02"}

价值是

{"before":null,"after":{"search_user_all_shards.Value":{"id":20,"name":{"string":"oliver"},"email":{"string":"bob@abc.com"},"department":{"string":"sales"},"modified":"2021-01-06T09:27:02Z"}},"source":{"version":"1.0.2.Final","connector":"mysql","name":"reviewDB","ts_ms":1609925222000,"snapshot":{"string":"false"},"db":"review","table":{"string":"search_user_02"},"server_id":1,"gtid":null,"file":"binlog.000002","pos":14630,"row":0,"thread":{"long":13},"query":null},"op":"c","ts_ms":{"long":1609925222956}}

作为接收器 hdfsSinkConnector 的一部分,我需要获取消息密钥 "__PKtableowner":"reviewDB.review.search_user_02 作为 hdfs 或 hive 中列或字段的一部分。

我发现的唯一 SMT 是 ValueToKey,但它似乎不适合我的用例,因为它是从值而不是从消息键中获取的。我已经尝试过(InsertField、CreateKey、ExtractField 等)您可以在此处找到的几乎所有转换,但没有运气。 https://docs.confluent.io/platform/current/connect/transforms/valuetokey.html

我正在寻找一种 KeyToValue 类型的 SMT,或者是否有其他解决方法。

以下是我的源和接收器配置。 来源:

{
  "name": "REVIEW__MYSQL__search_user__source",
  "config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "tasks.max": "1",
    "database.history.kafka.topic": "review.search_user_logs",
    "database.history.consumer.max.block.ms": "3000",
    "include.schema.changes": "false",
    "database.history.consumer.session.timeout.ms": "30000",
    "database.history.kafka.consumer.group": "compose-connect-group",
    "snapshot.new.tables": "parallel",
    "database.history.kafka.sasl.mechanism": "GSSAPI",
    "database.whitelist": "review",
    "database.history.producer.sasl.mechanism": "GSSAPI",
    "database.user": "root",
    "database.history.kafka.bootstrap.servers": "kafka:9092",
    "time.precision.mode": "connect",
    "database.server.name": "reviewDB",
    "database.port": "3306",
    "database.history.consumer.heartbeat.interval.ms": "1000",
    "min.row.count.to.stream.results": "0",
    "database.hostname": "mysql",
    "database.password": "example",
    "database.history.consumer.sasl.mechanism": "GSSAPI",
    "snapshot.mode": "when_needed",
    "table.whitelist": "review.search_user_(.*)",
    "transforms": "Reroute",
    "transforms.Reroute.type": "io.debezium.transforms.ByLogicalTableRouter",
    "transforms.Reroute.topic.regex": "reviewDB.review.search_user_(.*)",
    "transforms.Reroute.topic.replacement": "search_user_all_shards",
    "transforms.Reroute.key.field.name": "__PKtableowner"
  }
}

水槽

{ "name": "REVIEW__MYSQL__search_user__sink",
  "config":
  {
      "connector.class": "io.confluent.connect.hdfs.HdfsSinkConnector",
      "topics.dir": "/_incr_files",
      "flush.size": 1,
      "tasks.max": 1,
      "timezone": "UTC",
      "rotate.interval.ms": 5000,
      "locale": "en",
      "hadoop.home": "/etc/hadoop",
      "logs.dir": "/_incr_files_wal",
      "hive.integration": "false",
      "partition.duration.ms": "20000",
      "hadoop.conf.dir": "/etc/hadoop",
      "topics": "search_user_all_shards",
      "hdfs.url": "hdfs://namenode:9000",
      "transforms": "unwrap,insertTopicOffset,insertTimeStamp",
      "transforms.insertTimeStamp.type": "org.apache.kafka.connect.transforms.InsertField$Value",
      "transforms.unwrap.drop.tombstones": "true",
      "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
      "transforms.unwrap.delete.handling.mode": "rewrite",
      "transforms.insertTimeStamp.timestamp.field": "spdb_landing_timestamp",
      "transforms.insertTopicOffset.offset.field": "spdb_topic_offset",
      "transforms.insertTopicOffset.type": "org.apache.kafka.connect.transforms.InsertField$Value",
      "schema.compatibility": "NONE",
      "path.format": "'partition'=YYYY-MM-dd-HH",
      "partitioner.class": "io.confluent.connect.hdfs.partitioner.TimeBasedPartitioner"
  }
}

【问题讨论】:

    标签: apache-kafka hdfs debezium


    【解决方案1】:

    【讨论】:

      【解决方案2】:

      由于您的密钥是一个结构,我知道的最好的方法是这个 SMT,它有效地将密钥和值包装成一个新的嵌套值

      https://github.com/jcustenborder/kafka-connect-transform-archive

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2019-03-29
        • 1970-01-01
        • 2022-07-29
        • 2020-08-20
        • 2020-06-29
        • 2017-02-27
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多