如何通过Kafka Connect提取消息Key并追加至Value中
解决Kafka Connect将消息Key追加到Value的SMT配置问题
你当前的配置错误在于使用了InsertField$Value的静态字段插入参数,而非提取消息Key的参数。要实现把消息Key追加到Value中的需求,只需修改SMT的配置项即可:
修正后的连接器配置
{ "name": "aaa", "config": { "name": "aaa", "connector.class": "io.aiven.kafka.connect.opensearch.OpensearchSinkConnector", "tasks.max": "1", "topics": "test_x", "transforms": "InsertKeyToValue", "transforms.InsertKeyToValue.type": "org.apache.kafka.connect.transforms.InsertField$Value", "transforms.InsertKeyToValue.key.field": "key", "key.ignore": "false", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false", "key.converter.schemas.enable": "false", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "connection.url": "https://localhost:9200", "connection.username": "test", "connection.password": "test", "batch.size": "10000", "linger.ms": "1000", "errors.tolerance": "all", "errors.log.include.messages": "true", "errors.log.enable": "true" } }
关键配置说明
- 移除原配置中
static.field和static.value这两个插入固定值的参数,替换为key.field": "key":该参数会指示InsertField$ValueSMT提取消息Key的值,插入到Value的key字段中。 - 保留
key.converter": "org.apache.kafka.connect.storage.StringConverter":你的消息Key是字符串类型,该转换器可正确解析Key的值。 key.ignore": "false"确保连接器不会忽略消息Key,保证SMT能获取到Key的值。
配置生效后,发送的消息Value {"name":"ali"}和Key Person会被处理成{"name":"ali","key":"Person"}并写入OpenSearch。
内容的提问来源于stack exchange,提问作者Alihossein shahabi
相关产品推荐
相关产品推荐

