如何在Source Connector(如SpoolDirCsvSourceConnector)中捕获错误记录至指定Topic?
SpoolDirCsvSourceConnector的错误记录死信队列配置方法
是的,SpoolDirCsvSourceConnector这类Kafka Connect Source Connector,同样支持通过配置Kafka Connect核心的错误处理属性,将解析或处理过程中引发异常的错误记录发送到指定的死信队列Topic,配置逻辑和你提到的Jdbc Sink Connector基本一致。
你需要在Source Connector的配置中添加以下核心属性:
"errors.tolerance": "all", "errors.deadletterqueue.topic.name": "dead_topic", "errors.deadletterqueue.topic.replication.factor": 1, "errors.deadletterqueue.context.headers.enable": "true" // 可选,用于在死信记录中附加错误上下文信息
各属性说明:
errors.tolerance: 设置为all时,连接器会跳过错误记录并持续运行,同时将错误记录转发到死信队列;若设置为none,连接器遇到错误会直接停止。errors.deadletterqueue.topic.name: 指定接收错误记录的死信队列Topic名称。errors.deadletterqueue.topic.replication.factor: 指定死信队列Topic的副本因子,需匹配你的Kafka集群配置。errors.deadletterqueue.context.headers.enable: 可选配置,开启后会在死信记录的Headers中添加异常详情、源记录内容等信息,方便后续问题排查。
这些是Kafka Connect框架层面的通用错误处理配置,适用于所有遵循Kafka Connect规范的Source和Sink Connector,自然也包括SpoolDirCsvSourceConnector。
内容的提问来源于stack exchange,提问作者Muhamed Risvan M S
相关产品推荐
相关产品推荐

