概述
Datadog Logs Sink Connector 从 Kafka Topic 消费日志记录,通过 HTTPS 将日志发送到 Datadog,适用于将应用日志、服务运行事件等汇集到统一的日志检索链路。官方项目提供该 Connector。 每条非空记录的 value 转换为 JSON 后放入日志的message 字段,外层附带 ddsource=kafka-connect 和 Topic 标签,可统一补充服务、主机及环境标签。Kafka key 不写入日志正文;value 中的字段不会自动提升为 Datadog 的外层元数据。Connector 使用 Datadog Logs API 发送日志,HTTP 接收成功不等于日志已经索引或可查询。
前置条件
- 准备目标 Datadog 组织的有效 API key,并确认组织所属的 Datadog Site,避免将日志发送到其他站点。
授权许可
使用 Apache License 2.0。快速开始
提前准备 Connect Cluster、Kafka、日志 Topic 和目标 Datadog 组织,确认网络连通及访问权限;集群与 Connector 的准备和管理操作参见管理 Connector。以下配置适用于 key 为字符串或空值、value 为 UTF-8 日志文本的 Topic。datadoghq.com 或 datadoghq.eu,再应用配置。真实密钥应由受控的凭证管理方式提供,不要提交到版本库。字符串 value 会成为 JSON 字符串形式的 message;即使文本内容是 JSON,也不会自动解析为对象。接入后发送少量日志,在 Datadog 按 source:kafka-connect 和 topic:<logs-topic> 查找,并结合消费进度确认链路;无需等待积满一批才发送。
配置
Connector 与订阅
connector.class
指定 Connector 实现类。
- 类型:
string - 默认值:无
- 重要级别:高
- 必填:是
- 有效值 / 注意事项:使用
com.datadoghq.connect.logs.DatadogLogsSinkConnector。
topics
指定消费的日志 Topic。
- 类型:
list - 默认值:空列表
- 重要级别:高
- 有效值 / 注意事项:逗号分隔;与
topics.regex必须且只能设置一个非空订阅,不能包含 DLQ Topic。
topics.regex
按名称正则表达式订阅日志 Topic。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:使用合法的 Java 正则表达式;与
topics互斥,不得匹配 DLQ Topic。
tasks.max
设置 Task 数量上限。
- 类型:
int - 默认值:
1 - 重要级别:高
- 有效值 / 注意事项:至少为
1;有效消费并行度受 Topic 分区数和任务分配约束,不代表吞吐量按比例增长。
tasks.max.enforce
控制是否强制检查生成的 Task 数量上限。
- 类型:
boolean - 默认值:
true - 重要级别:低
- 已弃用:是
- 替代项:无独立替代参数,保持生成的任务数量符合
tasks.max。 - 有效值 / 注意事项:
true或false;不建议关闭数量上限检查。
Datadog 认证与目标站点
datadog.api_key
用于 Datadog 日志接收接口鉴权。
- 类型:
password - 默认值:无
- 重要级别:高
- 必填:是
- 有效值 / 注意事项:目标组织的有效 API key;空字符串不是有效凭证。
password类型不代表端到端凭证加密。
datadog.site
指定 Datadog Site,以确定日志接收主机。
- 类型:
string - 默认值:
null - 重要级别:中
- 有效值 / 注意事项:填写 Site 域名,例如
datadoghq.eu,不是控制台完整 URL。非空datadog.url优先;两项均未提供或为空时,运行时回退到http-intake.logs.datadoghq.com。该回退不是本项的默认值。
datadog.url
覆盖日志接收主机。
- 类型:
string - 默认值:
null - 重要级别:中
- 有效值 / 注意事项:填写主机和可选端口,例如
http-intake.logs.datadoghq.com:443,不要包含https://或路径;Connector 自动拼接 HTTPS 和/api/v2/logs。非空时覆盖datadog.site,应确认目标可信且支持该接口。
日志元数据
datadog.tags
为所有发送的日志添加公共标签。
- 类型:
list - 默认值:
null - 重要级别:中
- 有效值 / 注意事项:逗号分隔的标签列表;Connector 还会添加
topic:<topic-name>标签。
datadog.service
为日志设置统一的外层服务名。
- 类型:
string - 默认值:
null - 重要级别:中
- 有效值 / 注意事项:未设置时省略
service;不从每条记录的 value 自动提取服务名。
datadog.hostname
为日志设置统一的外层主机名。
- 类型:
string - 默认值:
null - 重要级别:中
- 有效值 / 注意事项:未设置时省略
hostname;不自动使用 Kafka Broker 或 Worker 的主机名。
datadog.add_published_date
将 Kafka 记录时间戳附加到日志。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true时为具有非空记录时间戳的日志添加毫秒值published_date;不能据此认为 Datadog 自动将其作为标准事件时间。
datadog.parse_record_headers
将 Kafka Headers 附加为 kafkaheaders 对象。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true或false;启用前确认 Header 转换器与实际数据兼容,并检查 Header 是否包含敏感信息。不保证任意 Header 结构都能转换。
代理
datadog.proxy.url
设置 HTTP 代理主机。
- 类型:
string - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:填写代理主机名,不是完整 URL;启用时须同时设置
datadog.proxy.port。没有公开的代理认证参数。
datadog.proxy.port
设置 HTTP 代理端口。
- 类型:
int - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:与非空代理主机配套使用,填写实际有效端口;没有默认代理端口。
发送重试
datadog.retry.max
设置连续发送失败时的 Connector 重试预算。
- 类型:
int - 默认值:
5 - 重要级别:低
- 有效值 / 注意事项:使用非负整数;
0表示不进行 Connector 重试。预算针对失败的put,成功后重置;不是严格的底层网络请求次数上限。重试耗尽后 Task 失败。
datadog.retry.backoff_ms
设置发送重试的基础等待时间,单位毫秒。
- 类型:
int - 默认值:
3000 - 重要级别:低
- 有效值 / 注意事项:使用正整数;后续等待会进行退避及随机化,不是精确重发定时器或总重试时限。不同于框架
errors.retry.*。
数据转换与变换
key.converter
指定 Kafka key 的转换器。
- 类型:
class - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:未设置时继承 Worker 转换器;虽然 key 不写入 Datadog,框架仍会转换 key,因此必须匹配上游编码。
value.converter
指定 Kafka value 的转换器。
- 类型:
class - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:未设置时继承 Worker 转换器;需与上游序列化格式匹配,额外参数属于所选转换器。转换后的 value 再转为 JSON 放入
message,不是直接将原始字节发送到 Datadog。
header.converter
指定 Kafka Headers 的转换器。
- 类型:
class - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:未设置时继承 Worker Header 转换器;公开 Headers 到日志时需检查类型兼容性。
transforms
声明按顺序执行的单消息变换(SMT)别名。
- 类型:
list - 默认值:空列表
- 重要级别:低
- 有效值 / 注意事项:别名不得重复,每个别名需要对应的变换类型;具体变换参数由所选插件定义。
transforms.<alias>.type
指定某个 SMT 别名的实现类。
- 类型:
class - 默认值:无
- 重要级别:高
- 必填:
required(配置该 SMT 别名时) - 有效值 / 注意事项:将
<alias>替换为transforms中的别名,使用可实例化的 Transformation 实现类。
predicates
声明用于条件执行 SMT 的谓词别名。
- 类型:
list - 默认值:空列表
- 重要级别:低
- 有效值 / 注意事项:别名不得重复;各谓词的专用参数由插件定义。
predicates.<alias>.type
指定某个谓词别名的实现类。
- 类型:
class - 默认值:无
- 重要级别:高
- 必填:
required(配置该谓词别名时) - 有效值 / 注意事项:将
<alias>替换为predicates中的别名,使用可实例化的 Predicate 实现类。
transforms.<alias>.predicate
为某个 SMT 选择执行条件。
- 类型:
string - 默认值:
null - 重要级别:中
- 有效值 / 注意事项:引用已配置的谓词别名;
<alias>是 SMT 别名。
transforms.<alias>.negate
控制是否反转 SMT 的谓词结果。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:显式设置本项时,必须同时显式设置对应的
predicate,即使本项值为false。
外部凭证配置更新
config.action.reload
控制 Config Provider 值变化时的处理方式。
- 类型:
string - 默认值:
restart - 重要级别:低
- 有效值 / 注意事项:
none或restart;Provider 的注册是 Worker 配置,本项本身不会启用凭证管理,也不保证任意密钥变化都能被自动检测。
框架错误处理
errors.retry.timeout
设置框架错误处理的总重试时长,单位毫秒。
- 类型:
long - 默认值:
0 - 重要级别:中
- 有效值 / 注意事项:
0不重试,-1无限重试;作用于框架支持的转换和 SMT 等阶段,不能替代 Datadog HTTP 发送重试。
errors.retry.delay.max.ms
设置框架重试的最大等待间隔,单位毫秒。
- 类型:
long - 默认值:
60000 - 重要级别:中
- 有效值 / 注意事项:不是
datadog.retry.backoff_ms,不控制本插件 HTTP 发送的退避。
errors.tolerance
控制框架支持的记录处理错误是否允许跳过。
- 类型:
string - 默认值:
none - 重要级别:中
- 有效值 / 注意事项:
none或all;all可容忍支持范围内的转换或 SMT 错误,但不自动跳过put内的 HTTP 发送或日志序列化失败。跳过意味着该记录未发送到 Datadog。
errors.log.enable
启用框架失败记录日志。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true或false;与 Connector 自身的异常日志独立,不改变失败处理方式。
errors.log.include.messages
控制框架错误日志是否包含记录的详细上下文。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:启用框架错误日志后生效;Sink 上下文包括 Topic、分区、Offset 和时间戳,应限制日志访问权限。
errors.deadletterqueue.topic.name
设置框架死信队列(DLQ)Topic。
- 类型:
string - 默认值:空字符串
- 重要级别:中
- 有效值 / 注意事项:空值禁用;与
errors.tolerance=all配合可保留框架支持范围内的失败记录。不得被本 Connector 订阅;不自动捕获本插件put内的 HTTP 发送失败。
errors.deadletterqueue.topic.replication.factor
设置自动创建 DLQ Topic 时的副本数。
- 类型:
short - 默认值:
3 - 重要级别:中
- 有效值 / 注意事项:用于 DLQ Topic 不存在时的创建,应符合 Kafka 集群可用 Broker 数量。
errors.deadletterqueue.context.headers.enable
为 DLQ 记录添加错误上下文 Headers。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true时添加__connect.errors.*上下文;需要已启用 DLQ,不扩展其可捕获的错误范围。
最佳实践
接入时补充统一检索标签
适用业务场景:日志能够送达后,需要在同一 Datadog 组织中区分服务和环境,便于后续检索与运维。先按共享同一服务身份的日志流配置,再接入更多来源,避免混合服务被标成同一服务。 配置示例:在快速开始配置中添加以下项;占位符分别代表这一日志流的服务名和部署环境。service 覆盖配置值;公共标签不应包含密钥或个人信息。
运行中为短暂发送失败留出恢复时间
适用业务场景:运行中遇到短暂的接收端故障,需要允许重试并观察积压,而不是立即停止任务。提高重试预算前,应先确认失败不是无效密钥、错误站点或稳定复现的数据问题。 配置示例:在快速开始配置中添加以下项,覆盖 Connector 的默认重试预算和基础等待时间。扩展时按分区增加消费并行度
适用业务场景:持续运行后出现消费积压,输入 Topic 有多个分区且目标端仍有接收余量,需要增加消费任务来分担处理。先观察延迟与错误,再逐步扩展,避免将下游瓶颈转为更多失败请求。 配置示例:对至少有两个可分配分区的输入,将快速开始配置的任务上限设置为以下值。监控
监控内容
关注 Kafka Connect 集群健康、Connector 和 Task 状态、输入及处理吞吐、消费积压与端到端延迟、Offset 提交进度及失败、错误与重试频率,以及 Worker JVM 的内存、GC 和线程信号;只有启用相应错误处理时才关注 DLQ 活动。持续积压或频繁重试时结合任务异常检查认证、站点和下游响应,不能仅凭RUNNING、输入记录计数或 HTTP 成功判断日志已被索引。
导入 Grafana 大盘
下载 AutoMQ Connect Cluster Dashboard,确认 Connect 指标已采集到 Prometheus 兼容数据源,且集群与实例等标签能够匹配大盘查询,再在 Grafana 导入 JSON 并选择对应数据源。限制条件
- 不提供恰好一次写入或幂等去重;部分请求成功后的重试、响应丢失及发送后尚未提交 Offset 的重启均可能产生重复日志。
- value 为
null的记录会跳过,不向 Datadog 发送删除指令,因此已消费的 Kafka 记录数不等于实际发送日志数。 - 不保证跨 Topic、分区或 Task 的全局顺序,也不保证 Datadog 的显示顺序。
- HTTP 成功是发送接受边界,不保证逐条日志已持久化、索引或查询可见;
flush也不会检查 Datadog 下游状态。 - 日志仍受 Datadog 接口及组织接收策略约束;Connector 会按内部约 4.5 MB 的未压缩 JSON 批次阈值分批,超过该阈值且无法单独放入空批次的单条序列化日志会被跳过。自动分批不代表任意大小的单条日志都能被 Datadog 接受,具体约束参见 Logs API。
常见问题
Task 为 RUNNING,但在 Datadog 中搜不到日志?
先检查输入 Topic 是否有非空 value、消费进度是否推进,以及 Task 日志中是否有发送或转换异常。确认 API key 与 Site 属于目标组织,搜索时间范围合适,并使用source:kafka-connect 和对应 Topic 标签检查。若 HTTP 已被接受但仍不可搜索,继续检查 Datadog 的日志处理、索引和排除策略;不要用 HTTP 成功代替查询可见性判断。
HTTP 发送反复失败,设置 DLQ 后仍然停止?
框架 DLQ 主要用于支持范围内的转换和 SMT 错误,不自动接收该插件put 内的 HTTP 失败。检查完整异常和安全的响应上下文:认证错误需修复 API key 与 Site,临时接收失败可等待 Connector 重试;持续失败或重试耗尽后应修复原因并恢复 Task。该插件不会按 HTTP 状态区分永久与临时错误,也不按 Retry-After 指定的时间调度。
出现重复日志,或重启后部分日志再次发送?
日志接收与 Kafka 消费 Offset 提交不是同一事务。发送成功但提交前退出,或整批处理中后续请求失败,都可能导致重发。检查重试、任务重启和 Offset 提交异常,将重复纳入下游处理设计;不要把增加重试预算理解为无重复保障。原始日志是 JSON,但字段仍在 message 中?
快速开始使用字符串转换器,JSON 文本仍是字符串,而不是自动解析的对象。需要结构化 Connect value 时,应选择与上游实际序列化格式匹配的转换器及其参数,再确认输出结构;即便 value 已是对象,其字段也留在message 内,不会自动变成外层 service 或 hostname。
配置接收地址后提示 URL 错误?
检查datadog.url 是否误填了完整 HTTPS 地址或 /api/v2/logs 路径。该项只接收主机和可选端口,由 Connector 拼接协议与路径;非空时还会覆盖 datadog.site。正常站点接入使用 Site 即可,移除不需要的 URL 覆盖。