Aerospike Kafka Outbound Connector Avro格式下map类型报错求助
Aerospike Kafka Outbound Connector Avro格式Map类型支持问题
我正尝试配置Aerospike Kafka Outbound Connector(源连接器)用于变更数据捕获(CDC)。按照官方文档使用Kafka Avro格式发送消息时,连接器抛出如下错误:
aerospike-kafka_connector-1 | 2022-12-22 20:27:17.607 GMT INFO metrics-ticker - requests-total: rate(per second) mean=8.05313731102431, m1=9.48698171570335, m5=2.480116641993411, m15=0.8667674157832074 aerospike-kafka_connector-1 | 2022-12-22 20:27:17.613 GMT INFO metrics-ticker - requests-total: duration(ms) min=1.441459, max=3101.488585, mean=15.582822432504553, stddev=149.48409869767494, median=4.713083, p75=7.851875, p95=17.496458, p98=28.421125, p99=85.418959, p999=3090.952252 aerospike-kafka_connector-1 | 2022-12-22 20:27:17.624 GMT ERROR metrics-ticker - **java.lang.Exception - Map type not allowed, has to be record type**: count=184
我的数据包含一个map字段,对应的Avro schema如下:
{ "name": "mydata", "type": "record", "fields": [ { "name": "metadata", "type": { "name": "com.aerospike.metadata", "type": "record", "fields": [ { "name": "namespace", "type": "string" }, { "name": "set", "type": [ "null", "string" ], "default": null }, { "name": "userKey", "type": [ "null", "long", "double", "bytes", "string" ], "default": null }, { "name": "digest", "type": "bytes" }, { "name": "msg", "type": "string" }, { "name": "durable", "type": [ "null", "boolean" ], "default": null }, { "name": "gen", "type": [ "null", "int" ], "default": null }, { "name": "exp", "type": [ "null", "int" ], "default": null }, { "name": "lut", "type": [ "null", "long" ], "default": null } ] } }, { "name": "test", "type": "string" }, { "name": "testmap", "type": { "type": "map", "values": "string", "default": {} } } ] }
是否有人成功实现过包含map字段的Avro格式消息发送场景?
补充说明:
看起来Aerospike连接器不支持map类型。我开启了详细日志,得到完整错误栈:
aerospike-kafka_connector-1 | 2023-01-06 18:41:31.098 GMT ERROR ErrorRegistry - Error stack trace aerospike-kafka_connector-1 | java.lang.Exception: Map type not allowed, has to be record type aerospike-kafka_connector-1 | at com.aerospike.connect.kafka.outbound.parser.KafkaAvroOutboundRecordGenerator$Companion.assertOnlyValidTypes(KafkaAvroOutboundRecordGenerator.kt:155) aerospike-kafka_connector-1 | at com.aerospike.connect.kafka.outbound.parser.KafkaAvroOutboundRecordGenerator$Companion.assertOnlyValidTypes(KafkaAvroOutboundRecordGenerator.kt:164) aerospike-kafka_connector-1 | at com.aerospike.connect.kafka.outbound.parser.KafkaAvroOutboundRecordGenerator$Companion.assertSchemaValid(KafkaAvroOutboundRecordGenerator.kt:101) aerospike-kafka_connector-1 | at com.aerospike.connect.kafka.outbound.parser.KafkaAvroOutboundRecordGenerator.<init>(KafkaAvroOutboundRecordGenerator.kt:180) aerospike-kafka_connector-1 | at com.aerospike.connect.kafka.outbound.inject.KafkaOutboundGuiceModule.getKafkaAvroStreamingRecordParser(KafkaOutboundGuiceModule.kt:59) aerospike-kafka_connector-1 | at com.aerospike.connect.kafka.outbound.inject.KafkaOutboundGuiceModule.access$getKafkaAvroStreamingRecordParser(KafkaOutboundGuiceModule.kt:29) aerospike-kafka_connector-1 | at com.aerospike.connect.kafka.outbound.inject.KafkaOutboundGuiceModule$bindKafkaAvroParserFactory$1.get(KafkaOutboundGuiceModule.kt:48) aerospike-kafka_connector-1 | at com.aerospike.connect.outbound.converter.XdrExchangeConverter.getInbuiltRecordFormatter(XdrExchangeConverter.kt:422) aerospike-kafka_connector-1 | at com.aerospike.connect.outbound.converter.XdrExchangeConverter.access$getInbuiltRecordFormatter(XdrExchangeConverter.kt:75) aerospike-kafka_connector-1 | at com.aerospike.connect.outbound.converter.XdrExchangeConverter$RouterAndInbuiltFormatter.<init>(XdrExchangeConverter.kt:285) aerospike-kafka_connector-1 | at com.aerospike.connect.outbound.converter.XdrExchangeConverter.processXdrRecord(XdrExchangeConverter.kt:192) aerospike-kafka_connector-1 | at com.aerospike.connect.outbound.converter.XdrExchangeConverter.parse(XdrExchangeConverter.kt:134) aerospike-kafka_connector-1 | at com.aerospike.connect.outbound.OutboundBridge$processAsync$1.invokeSuspend(OutboundBridge.kt:182) aerospike-kafka_connector-1 | at kotlin.coroutines.jvm.internal.BaseContinuationImpl.resumeWith(ContinuationImpl.kt:33) aerospike-kafka_connector-1 | at kotlinx.coroutines.DispatchedTask.run(DispatchedTask.kt:106) aerospike-kafka_connector-1 | at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) aerospike-kafka_connector-1 | at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) aerospike-kafka_connector-1 | at java.base/java.lang.Thread.run(Thread.java:829)
但该限制并未在官方文档中提及。
内容的提问来源于stack exchange,提问作者arunvelsriram
相关产品推荐
相关产品推荐

