Skip to main content

概述

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。
替换日志 Topic、API key 和 Site 域名,例如 datadoghq.comdatadoghq.eu,再应用配置。真实密钥应由受控的凭证管理方式提供,不要提交到版本库。字符串 value 会成为 JSON 字符串形式的 message;即使文本内容是 JSON,也不会自动解析为对象。接入后发送少量日志,在 Datadog 按 source:kafka-connecttopic:<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
  • 有效值 / 注意事项truefalse;不建议关闭数量上限检查。

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
  • 重要级别:中
  • 有效值 / 注意事项truefalse;启用前确认 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
  • 重要级别:低
  • 有效值 / 注意事项nonerestart;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
  • 重要级别:中
  • 有效值 / 注意事项noneallall 可容忍支持范围内的转换或 SMT 错误,但不自动跳过 put 内的 HTTP 发送或日志序列化失败。跳过意味着该记录未发送到 Datadog。

errors.log.enable

启用框架失败记录日志。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项truefalse;与 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 组织中区分服务和环境,便于后续检索与运维。先按共享同一服务身份的日志流配置,再接入更多来源,避免混合服务被标成同一服务。 配置示例:在快速开始配置中添加以下项;占位符分别代表这一日志流的服务名和部署环境。
关键说明:这些值统一应用到该 Connector 的日志外层,Topic 标签仍自动添加。不同服务可使用各自 Connector,不要期待 value 内的 service 覆盖配置值;公共标签不应包含密钥或个人信息。

运行中为短暂发送失败留出恢复时间

适用业务场景:运行中遇到短暂的接收端故障,需要允许重试并观察积压,而不是立即停止任务。提高重试预算前,应先确认失败不是无效密钥、错误站点或稳定复现的数据问题。 配置示例:在快速开始配置中添加以下项,覆盖 Connector 的默认重试预算和基础等待时间。
关键说明:示例值不是通用最优值,应依据可容忍的积压和恢复时间调整。同步发送与等待会形成背压;成功后预算重置,耗尽后需修复原因并恢复 Task。重试可能重发已被接收的日志,框架 DLQ 不能代替此处的失败恢复;不要将重复日志当成恰好一次写入。

扩展时按分区增加消费并行度

适用业务场景:持续运行后出现消费积压,输入 Topic 有多个分区且目标端仍有接收余量,需要增加消费任务来分担处理。先观察延迟与错误,再逐步扩展,避免将下游瓶颈转为更多失败请求。 配置示例:对至少有两个可分配分区的输入,将快速开始配置的任务上限设置为以下值。
关键说明:任务分配由 Connect 管理,实际有效并行度受分区数限制。该值只是扩展起点,不承诺线性提速或全局顺序;增加任务后需同时观察 Datadog 接收错误、吞吐及消费积压。

监控

监控内容

关注 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 内,不会自动变成外层 servicehostname

配置接收地址后提示 URL 错误?

检查 datadog.url 是否误填了完整 HTTPS 地址或 /api/v2/logs 路径。该项只接收主机和可选端口,由 Connector 拼接协议与路径;非空时还会覆盖 datadog.site。正常站点接入使用 Site 即可,移除不需要的 URL 覆盖。

排障日志中出现业务内容,如何控制泄露风险?

不要在生产环境长期启用 TRACE;请求正文、响应和部分超大日志的预览可能进入运行日志。限制日志级别、访问权限与保留时间,避免将敏感 Header 暴露到 Datadog,并确认自定义接收主机可信。共享诊断材料前应脱敏,但保留异常堆栈和安全的 Topic、分区、Offset 等定位信息。