> ## Documentation Index
> Fetch the complete documentation index at: https://docs.automq.com/llms.txt
> Use this file to discover all available pages before exploring further.

# Datadog Logs Sink Connector

> 介绍如何在 AutoMQ Connect 中配置和运行 Datadog Logs Sink Connector，包括前置条件、配置、监控和故障排查。

## 概述

Datadog Logs Sink Connector 从 Kafka Topic 消费日志记录，通过 HTTPS 将日志发送到 Datadog，适用于将应用日志、服务运行事件等汇集到统一的日志检索链路。[官方项目](https://github.com/DataDog/datadog-kafka-connect-logs)提供该 Connector。

每条非空记录的 value 转换为 JSON 后放入日志的 `message` 字段，外层附带 `ddsource=kafka-connect` 和 Topic 标签，可统一补充服务、主机及环境标签。Kafka key 不写入日志正文；value 中的字段不会自动提升为 Datadog 的外层元数据。Connector 使用 [Datadog Logs API](https://docs.datadoghq.com/api/latest/logs/) 发送日志，HTTP 接收成功不等于日志已经索引或可查询。

## 前置条件

* 准备目标 Datadog 组织的有效 API key，并确认组织所属的 [Datadog Site](https://docs.datadoghq.com/getting_started/site/)，避免将日志发送到其他站点。

## 授权许可

使用 Apache License 2.0。

## 快速开始

提前准备 Connect Cluster、Kafka、日志 Topic 和目标 Datadog 组织，确认网络连通及访问权限；集群与 Connector 的准备和管理操作参见[管理 Connector](../manage-connectors)。以下配置适用于 key 为字符串或空值、value 为 UTF-8 日志文本的 Topic。

```properties theme={null}
connector.class=com.datadoghq.connect.logs.DatadogLogsSinkConnector
topics=<logs-topic>
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.storage.StringConverter
datadog.api_key=<datadog-api-key>
datadog.site=<datadog-site>
```

替换日志 Topic、API key 和 Site 域名，例如 `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 组织中区分服务和环境，便于后续检索与运维。先按共享同一服务身份的日志流配置，再接入更多来源，避免混合服务被标成同一服务。

**配置示例**：在快速开始配置中添加以下项；占位符分别代表这一日志流的服务名和部署环境。

```properties theme={null}
datadog.service=<service-name>
datadog.tags=env:<environment>
```

**关键说明**：这些值统一应用到该 Connector 的日志外层，Topic 标签仍自动添加。不同服务可使用各自 Connector，不要期待 value 内的 `service` 覆盖配置值；公共标签不应包含密钥或个人信息。

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

**适用业务场景**：运行中遇到短暂的接收端故障，需要允许重试并观察积压，而不是立即停止任务。提高重试预算前，应先确认失败不是无效密钥、错误站点或稳定复现的数据问题。

**配置示例**：在快速开始配置中添加以下项，覆盖 Connector 的默认重试预算和基础等待时间。

```properties theme={null}
datadog.retry.max=8
datadog.retry.backoff_ms=5000
```

**关键说明**：示例值不是通用最优值，应依据可容忍的积压和恢复时间调整。同步发送与等待会形成背压；成功后预算重置，耗尽后需修复原因并恢复 Task。重试可能重发已被接收的日志，框架 DLQ 不能代替此处的失败恢复；不要将重复日志当成恰好一次写入。

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

**适用业务场景**：持续运行后出现消费积压，输入 Topic 有多个分区且目标端仍有接收余量，需要增加消费任务来分担处理。先观察延迟与错误，再逐步扩展，避免将下游瓶颈转为更多失败请求。

**配置示例**：对至少有两个可分配分区的输入，将快速开始配置的任务上限设置为以下值。

```properties theme={null}
tasks.max=2
```

**关键说明**：任务分配由 Connect 管理，实际有效并行度受分区数限制。该值只是扩展起点，不承诺线性提速或全局顺序；增加任务后需同时观察 Datadog 接收错误、吞吐及消费积压。

## 监控

### 监控内容

关注 Kafka Connect 集群健康、Connector 和 Task 状态、输入及处理吞吐、消费积压与端到端延迟、Offset 提交进度及失败、错误与重试频率，以及 Worker JVM 的内存、GC 和线程信号；只有启用相应错误处理时才关注 DLQ 活动。持续积压或频繁重试时结合任务异常检查认证、站点和下游响应，不能仅凭 `RUNNING`、输入记录计数或 HTTP 成功判断日志已被索引。

### 导入 Grafana 大盘

下载 [AutoMQ Connect Cluster Dashboard](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，确认 Connect 指标已采集到 Prometheus 兼容数据源，且集群与实例等标签能够匹配大盘查询，再在 Grafana 导入 JSON 并选择对应数据源。

## 限制条件

* 不提供恰好一次写入或幂等去重；部分请求成功后的重试、响应丢失及发送后尚未提交 Offset 的重启均可能产生重复日志。
* value 为 `null` 的记录会跳过，不向 Datadog 发送删除指令，因此已消费的 Kafka 记录数不等于实际发送日志数。
* 不保证跨 Topic、分区或 Task 的全局顺序，也不保证 Datadog 的显示顺序。
* HTTP 成功是发送接受边界，不保证逐条日志已持久化、索引或查询可见；`flush` 也不会检查 Datadog 下游状态。
* 日志仍受 Datadog 接口及组织接收策略约束；Connector 会按内部约 4.5 MB 的未压缩 JSON 批次阈值分批，超过该阈值且无法单独放入空批次的单条序列化日志会被跳过。自动分批不代表任意大小的单条日志都能被 Datadog 接受，具体约束参见 [Logs API](https://docs.datadoghq.com/api/latest/logs/)。

## 常见问题

### 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 覆盖。

### 排障日志中出现业务内容，如何控制泄露风险？

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