【问题标题】:Unable to parse a topic from Kafka Connect AVRO connector无法从 Kafka Connect AVRO 连接器解析主题
【发布时间】:2022-07-11 20:15:35
【问题描述】:

我正在使用 Neo4j 接收器连接器从 Kafka 主题读取数据并将其转储到 Neo4j 数据库。 Kafka 中可用的消息/数据是 AVRO 格式,因此我尝试使用 AVRO 转换器通过提供模式注册表详细信息来解析数据。但是在使用该消息时,我看到了一个 DataError 异常。

以下是我创建连接器的配置。

{
    "topics": "mytopic",
    "connector.class": "streams.kafka.connect.sink.Neo4jSinkConnector",
    "tasks.max":"1",
    "key.converter.schemas.enable":"true",
    "values.converter.schemas.enable":"true",
    "errors.retry.timeout": "-1",
    "errors.retry.delay.max.ms": "1000",
    "errors.tolerance": "none",
    "errors.deadletterqueue.topic.name": "deadletter-topic",
    "errors.deadletterqueue.topic.replication.factor":1,
    "errors.deadletterqueue.context.headers.enable":true,
    "key.converter":"org.apache.kafka.connect.storage.StringConverter",
    "key.converter.enhanced.avro.schema.support":true,
    "value.converter.enhanced.avro.schema.support":true,
    "value.converter":"io.confluent.connect.avro.AvroConverter",
    "value.converter.schema.registry.url":"https://schema-url/",
    "value.converter.basic.auth.credentials.source":"USER_INFO",
    "value.converter.basic.auth.user.info":"user:pass",
    "errors.log.enable": true,
    "schema.ignore":"false",
    "errors.log.include.messages": true,
    "neo4j.server.uri": "neo4j://my-ip:7687/neo4j",
    "neo4j.authentication.basic.username": "neo4j",
    "neo4j.authentication.basic.password": "neo4j",
    "neo4j.encryption.enabled": false,
    "neo4j.topic.cypher.mytopic": "MERGE (p:Loc_Con{name: event.geography.name})"
}

这是我得到的例外。

ErrorData(originalTopic=mytopic, timestamp=1652188554497, partition=0, offset=2140111, exception=org.apache.kafka.connect.errors.DataException: Exception thrown while processing field 'geography', key=9662840       , value=Struct{geography=Struct{geoId=43333,geoType=Business Defined Area,name=Norarea,status=Active,validFrom=Sat Apr 09 00:00:00 GMT 2012,validTo=Fri Dec 31 00:00:00 GMT 9999, executingClass=class streams.kafka.connect.sink.Neo4jSinkTask)

我想知道这里出了什么问题,我也尝试过使用 String 和 JSON Converter,但是解析也失败了。那么有没有办法解析数据呢?

【问题讨论】:

  • 这是完整的堆栈跟踪吗?
  • @OneCricketeer ,是的,它是完整的堆栈,我发现了这个问题。问题是我需要从收到的消息中提取字段 geogrpahy ,使用下面的转换来提取。 “transforms”:“ExtractField”,“transforms.ExtractField.type”:“org.apache.kafka.connect.transforms.ExtractField$Value”,“transforms.ExtractField.field”:“geography”
  • 当然@OneCricketeer

标签: apache-kafka neo4j apache-kafka-connect


【解决方案1】:

我发现了问题所在。我需要进行转换以提取名为 geography 的数据字段。它基本上从整个 JSON 中提取 geography 字段并将其分配回“事件”。

"transforms": "ExtractField", 
"transforms.ExtractField.type": "org.apache.kafka.connect.transforms.ExtractField$Value", 
"transforms.ExtractField.field": "geography"

【讨论】:

    猜你喜欢
    • 2020-04-29
    • 2021-09-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-09-03
    • 1970-01-01
    • 2019-07-03
    • 2018-01-12
    相关资源
    最近更新 更多