Kafka Connect对接Azure Data Lake Sink任务失败求助
问题概述
使用InstaClustr托管的Kafka Bootstrap Server和Kafka Connect,部署Axual的ADLS Gen2 Sink连接器后,启动出现错误Null scanresults or scanresults types are not supported;通过控制台生产者发送消息后,Azure Data Lake Storage Gen2中未生成任何内容,已提前创建配置指定的base目录。
连接器配置
{ "name": "my-connector", "config": { "connector.class": "io.axual.connect.plugins.adls.gen2.AdlsGen2SinkConnector", "tasks.max": "3", "topics": "test", "adls.endpoint": "https://accountname.dfs.core.windows.net", "adls.container.name": "test", "adls.auth.method": "AccountKey", "adls.account.name": "myexampleadlsaccount", "adls.account.key": "get-this-from-the-azure-storage-account-settings", "base.directory": "base" } }
完整错误栈信息
{"state":"FAILED","trace":"org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception. at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:609) at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:329) at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:232) at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:201) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:186) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:241) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: io.axual.connect.plugins.adls.gen2.exceptions.AdlsGen2Exception: io.axual.connect.plugins.adls.gen2.exceptions.AdlsGen2Exception: Null scanresults or scanresult types are not supported at io.axual.connect.plugins.adls.gen2.AdlsGen2SinkTask.put(AdlsGen2SinkTask.java:157) at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:581) ... 10 more Caused by: io.axual.connect.plugins.adls.gen2.exceptions.AdlsGen2Exception: Null scanresults or scanresult types are not supported at io.axual.connect.plugins.adls.gen2.avro.ContainerGenerator.determineInnerSchema(ContainerGenerator.java:118) at io.axual.connect.plugins.adls.gen2.avro.ContainerGenerator.createKafkaKeyValueContainer(ContainerGenerator.java:105) at io.axual.connect.plugins.adls.gen2.avro.ContainerGenerator.createContainer(ContainerGenerator.java:97) at io.axual.connect.plugins.adls.gen2.avro.AvroADLSFile.processRecords(AvroADLSFile.java:193) at io.axual.connect.plugins.adls.gen2.AdlsGen2SinkTask.lambda$put$4(AdlsGen2SinkTask.java:142) at java.base/java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:195) at java.base/java.util.HashMap$EntrySpliterator.forEachRemaining(HashMap.java:1764) at java.base/java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:484) at java.base/java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:474) at java.base/java.util.stream.ReduceOps$ReduceOp.evaluateSequential(ReduceOps.java:913) at java.base/java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234) at java.base/java.util.stream.ReferencePipeline.reduce(ReferencePipeline.java:558) at io.axual.connect.plugins.adls.gen2.AdlsGen2SinkTask.put(AdlsGen2SinkTask.java:143)\ ... 11 more","worker_id":"10.201.0.1:8083","generation":85}
排查与解决方法
1. 修正消息格式与转换器配置
错误栈显示问题出在Avro Schema解析环节(ContainerGenerator.determineInnerSchema),该连接器默认期望处理带Schema的Avro格式消息,若发送无Schema的原始字符串/JSON消息,会导致Schema扫描结果为空,触发报错。
- 若使用非Avro消息,需在连接器配置中添加转换器参数,以JSON消息为例:
"value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false", "key.converter": "org.apache.kafka.connect.storage.StringConverter" - 若使用Avro消息,需确保配置Schema Registry相关参数(如
value.converter.schema.registry.url),并保证生产者发送的消息附带合法Schema。
2. 校验连接器与Kafka版本兼容性
确认Axual连接器版本与InstaClustr托管的Kafka Connect版本匹配,避免因API差异导致Schema处理逻辑异常。
3. 统一ADLS配置参数
检查adls.endpoint中的账户名与adls.account.name是否一致(当前配置中两者分别为accountname和myexampleadlsaccount,需修正为同一账户名),同时验证ADLS账户密钥的有效性,确保连接器拥有容器及base目录的读写权限。
4. 启用调试日志定位问题
在Kafka Connect配置中提高Axual连接器的日志级别,将log4j.logger.io.axual.connect.plugins.adls.gen2设置为DEBUG,查看消息Schema扫描的详细过程,进一步定位空Schema的触发原因。
内容的提问来源于stack exchange,提问作者aaron

