【问题标题】:How to add column with the kafka message timestamp in kafka sink connector如何在 kafka sink 连接器中添加带有 kafka 消息时间戳的列
【发布时间】:2019-04-18 10:18:09
【问题描述】:

我正在使用属性/json 文件配置我的连接器,当它从源连接器读取消息但没有成功时,我试图添加一个包含 kafka 时间戳的时间戳列。

我尝试添加transforms,但它始终为空,并且我的接收器连接器“大查询”返回错误

更新表架构失败

我确实将这些配置放在 bigquery 连接器属性中

transforms=InsertField
transforms.InsertField.timestamp.field=fieldtime
transforms.InsertField.type=org.apache.kafka.connect.transforms.InsertField$Value

我的源 Config Sap 连接器

{
    "name": "sap",
    "config": {
        "connector.class": "com.sap.kafka.connect.source.hana.HANASourceConnector",
        "tasks.max": "10",
        "topics": "mytopic",
        "connection.url": "jdbc:sap://IP:30015/",
        "connection.user": "user",
        "connection.password": "pass",
        "group.id":"589f5ff5-1c43-46f4-bdd3-66884d61m185",
        "mytopic.table.name":                          "\"schema\".\"mytable\""  
       }
}

我的接收器连接器 BigQuery

name=bigconnect
connector.class=com.wepay.kafka.connect.bigquery.BigQuerySinkConnector
tasks.max=1

sanitizeTopics=true

autoCreateTables=true
autoUpdateSchemas=true

schemaRetriever=com.wepay.kafka.connect.bigquery.schemaregistry.schemaretriever.SchemaRegistrySchemaRetriever
schemaRegistryLocation=http://localhost:8081

bufferSize=100000
maxWriteSize=10000
tableWriteWait=1000

project=kafka-test-217517
topics=mytopic
datasets=.*=sap_dataset
keyfile=/opt/bgaccess.json
transforms=InsertField
transforms.InsertField.timestamp.field=fieldtime    
transforms.InsertField.type=org.apache.kafka.connect.transforms.InsertField$Value

【问题讨论】:

  • 您使用的是哪个接收器连接器?
  • bigquery 我已尝试使用 transforms.InsertSource.timestamp.field 但它给出了一个错误,即无法修改架构
  • 您能否分享您的完整 Kafka Connect 配置(工作程序和连接器)以及来自 Kafka Connect 工作程序的日志?
  • @RobinMoffatt 源是不同的多个数据库,目标连接器是 Bigquery 我用我的配置更新了我的答案
  • 这不是我问的:)你能分享你的 Kafka Connect 配置文件(worker 和 connector)的全部内容,以及来自 Kafka Connect worker 的日志吗?

标签: apache-kafka google-bigquery apache-kafka-connect


【解决方案1】:

老答案 我想我明白了背后的问题

首先,您不能在任何源连接器中使用转换 InsertField,因为 msg 的时间戳值是在写入主题时分配的,因此连接器不是已经知道的,
对于 JDBC 连接器,有这张票 https://github.com/confluentinc/kafka-connect-jdbc/issues/311

在 sap 源连接器中也无法正常工作。

第二个 BigQuery 连接器有一个错误,不允许使用 InsertField 将时间戳添加到此处提到的每个表

https://github.com/wepay/kafka-connect-bigquery/issues/125#issuecomment-439102994

因此,如果您想使用 bigquery 作为输出,现在唯一的解决方案是在加载 cink 连接器之前手动编辑每个表的架构以添加列

2018 年 12 月 3 日更新 始终在 SINK 连接器中添加消息时间戳的最终解决方案。假设您要将时间戳添加到接收器连接器的每个表中

在您的 SOURCE CONNECTOR 中放置此配置

"transforms":"InsertField"
"transforms.InsertField.timestamp.field":"fieldtime", 
"transforms.InsertField.type":"org.apache.kafka.connect.transforms.InsertField$Value"

这将为每个源表添加一个名为“fieldtime”的列名

在您的 SINK CONNECTOR 中放置这些配置

"transforms":"InsertField,DropField",
"transforms.DropField.type":"org.apache.kafka.connect.transforms.ReplaceField$Value",
"transforms.DropField.blacklist":"fieldtime",
"transforms.InsertSource.timestamp.field":"kafka_timestamp",
"transforms.InsertField.timestamp.field":"fieldtime",
"transforms.InsertField.type":"org.apache.kafka.connect.transforms.InsertField$Value"

这实际上将删除字段时间列并使用消息的时间戳再次添加它

此方案将自动添加正确值的列,无需任何添加操作

【讨论】:

    【解决方案2】:

    我猜您的错误来自 BigQuery,而不是 Kafka Connect。

    例如,在独立模式下启动 Connect Console Consumer,您会看到类似

    的消息

    Struct{...,fieldtime=Fri Nov 16 07:38:19 UTC 2018}


    connect-standalone ./connect-standalone.properties ./connect-console-sink.properties测试

    我有一个带有 Avro 数据的输入主题...相应地更新您自己的设置

    connect-standalone.properties

    bootstrap.servers=kafka:9092
    
    key.converter=io.confluent.connect.avro.AvroConverter
    key.converter.schema.registry.url=http://schema-registry:8081
    key.converter.schemas.enable=true
    
    value.converter=io.confluent.connect.avro.AvroConverter
    value.converter.schema.registry.url=http://schema-registry:8081
    value.converter.schemas.enable=true
    
    offset.storage.file.filename=/tmp/connect.offsets
    offset.flush.interval.ms=10000
    
    plugin.path=/usr/share/java
    

    connect-console-sink.properties

    name=local-console-sink
    connector.class=org.apache.kafka.connect.file.FileStreamSinkConnector
    tasks.max=1
    topics=input-topic
    
    transforms=InsertField
    transforms.InsertField.timestamp.field=fieldtime
    transforms.InsertField.type=org.apache.kafka.connect.transforms.InsertField$Value
    

    【讨论】:

    • 在我的情况下,源连接器和接收器连接器中存在问题。关于 BigQuery-Sink 实际上问题与无法很好地管理 transforms.InsertField.timestamp 的连接器有关。关于源连接器,如果您使用 transforms.InsertField.timestamp 这将始终为 NULL。所有主题都有 ROWTIME 列,但我无法使用它
    • KSQL 有一个“ROWTIME 列”...实际消息可能没有,因为 Kafka 消息中没有列之类的东西...话虽如此,它是 Java API 的 record.timestamp() ,这是timestamp.field=fieldtime 获得的字段...并且所有连接器都支持所有相同的转换,它并不特定于某个转换。无论如何,Confluent 有一个博客,包括使用 BigQuery confluent.io/blog/data-wrangling-apache-kafka-ksql
    猜你喜欢
    • 2021-01-26
    • 1970-01-01
    • 2013-09-16
    • 2019-08-18
    • 1970-01-01
    • 2021-12-22
    • 2022-07-28
    • 2020-01-24
    • 2017-11-12
    相关资源
    最近更新 更多