【问题标题】:How to have field type support for date in kafka如何在kafka中为日期提供字段类型支持
【发布时间】:2019-11-13 04:32:28
【问题描述】:

使用 debezium-mongodb-connector 我设法将我的集合推送到 kafka,我面临的唯一问题是我的一个集合中的字段 date 使用这种格式 2019-05-14T23:25:34.703+ 00:00,没有以相同的格式被推送到主题,而是我得到了类似 1560708085175 的内容。

这是我的 debezium 连接器命令connect-standalone /etc/kafka/connect-standalone.properties /etc/kafka/connect-mongodb-source.properties 这是我的 mongodb 集合示例。

{"_id":"5cdb4e6ed767ba70593e2aa8","sender":"5cdb43db4505956efc70ba03","receiver":"5cdb43db4505956efc70ba03","receiverWalletId":"5cdb43db4505956efc70ba04","status":"succes","type":"topup","amount":200000,"totalFee":0,"createdAt":"2019-05-14T23:25:34.703Z","updatedAt":"2019-05-14T23:25:35.132Z","__v":0,"details":"none."}

这是我的 kafka 主题示例。

{"schema":{"type":"struct","fields":[{"type":"string","optional":true,"field":"sender"},{"type":"string","optional":true,"field":"receiver"},{"type":"string","optional":true,"field":"receiverWalletId"},{"type":"string","optional":true,"field":"status"},{"type":"string","optional":true,"field":"type"},{"type":"int32","optional":true,"field":"amount"},{"type":"int32","optional":true,"field":"totalFee"},{"type":"int64","optional":true,"field":"createdAt"},{"type":"int64","optional":true,"field":"updatedAt"},{"type":"int32","optional":true,"field":"__v"},{"type":"string","optional":true,"field":"from"},{"type":"string","optional":true,"field":"orderId"},{"type":"string","optional":true,"field":"id"}],"optional":false,"name":"mongo_conn.digi.transactions"},"payload":{"sender":"5cef970ca2e9c273c655483","receiver":"5cef970ca2e9c27355c483","receiverWalletId":"5cef970ca2e9c27556c484","status":"pending","type":"topup","amount":6000,"totalFee":0,"createdAt":1560708024322,"updatedAt":1560708024753,"__v":0,"from":"smt","orderId":"d7a97581-9d18-79cd-8b09-16e400a43714","id":"5d0683b8be4af834abe3cf58"}}

这是我的 connect-mongodb-source.properties

name=mongodb-source-connector
connector.class=io.debezium.connector.mongodb.MongoDbConnector
mongodb.hosts=repracli/**.**.**.***27017
mongodb.name=mongo_conn
initial.sync.max.threads=1
tasks.max=1
transforms=unwrap
transforms.unwrap.type=io.debezium.connector.mongodb.transforms.UnwrapFromMongo$
transforms.unwrap.operation.header=true

【问题讨论】:

    标签: mongodb apache-kafka apache-kafka-connect debezium


    【解决方案1】:

    对于几个转换,您将需要以下内容:

    transforms=unwrap,convert1,convert2
    transforms.unwrap.type=io.debezium.connector.mongodb.transforms.UnwrapFromMongoDbEnvelope
    transforms.unwrap.operation.header=true
    transforms.convert1.type=org.apache.kafka.connect.transforms.TimestampConverter$Value
    transforms.convert1.target.type=string
    transforms.convert1.field=createdAt
    transforms.convert1.format=yyyy-MM-dd HH:mm:ss ZZZ
    transforms.convert2.type=org.apache.kafka.connect.transforms.TimestampConverter$Value
    transforms.convert2.target.type=string
    transforms.convert2.field= *updatedAt*
    transforms.convert2.format=yyyy-MM-dd HH:mm:ss ZZZ

    【讨论】:

      【解决方案2】:

      Debezium 以格式流式传输数据,因为它们存储在 oplog 中。日期看起来像自纪元以来以毫秒为单位的 unix 时间戳。

      您可以编写一个 SMT (https://cwiki.apache.org/confluence/display/KAFKA/KIP-66%3A+Single+Message+Transforms+for+Kafka+Connect) 来处理消息并将请求的字段转换为您喜欢的字符串表示形式。

      如果您查看org.bson.BsonDateTime,您会发现它确实是long 值。

      【讨论】:

      • 我已将这些行添加到我的 connect-mongodb-source.properties。 “transforms”:“TimestampConverter”,“transforms.TimestampConverter.field”:“createdAt”“transforms.TimestampConverter.type”:“org.apache.kafka.connect.transforms.TimestampConverter$Value”,“transforms.TimestampConverter.format” : "yyyy-MM-dd HH:mm:ss ZZZ" "transforms.TimestampConverter.target.type": "string" 它不起作用
      【解决方案3】:

      解决了

      name=mongodb-source-connector
      connector.class=io.debezium.connector.mongodb.MongoDbConnector
      mongodb.hosts=repracli/**.**.**.***:27017
      mongodb.name=mongo_conn
      initial.sync.max.threads=1
      tasks.max=1
      transforms=unwrap,convert,convert2,convert3,convert4
      transforms.unwrap.type=io.debezium.connector.mongodb.transforms.UnwrapFromMongoDbEnvelope
      transforms.unwrap.operation.header=true
      transforms.convert.type=org.apache.kafka.connect.transforms.TimestampConverter$Value
      transforms.convert.target.type=string
      transforms.convert.field=createdAt
      transforms.convert.format=yyyy-MM-dd HH:mm:ss ZZZ
      transforms.convert2.type=org.apache.kafka.connect.transforms.TimestampConverter$Value
      transforms.convert2.target.type=string
      transforms.convert2.field=updatedAt
      transforms.convert2.format=yyyy-MM-dd HH:mm:ss ZZZ
      transforms.convert3.type=org.apache.kafka.connect.transforms.TimestampConverter$Value
      transforms.convert3.target.type=string
      transforms.convert3.field=created_at
      transforms.convert3.format=yyyy-MM-dd HH:mm:ss ZZZ
      transforms.convert4.type=org.apache.kafka.connect.transforms.TimestampConverter$Value
      transforms.convert4.target.type=string
      transforms.convert4.field=updated_at
      transforms.convert4.format=yyyy-MM-dd HH:mm:ss ZZZ
      

      【讨论】:

        猜你喜欢
        • 2013-03-06
        • 1970-01-01
        • 2015-10-09
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多