概述
New Relic Sink Connector 从 Kafka Topic 消费记录,将记录转换为 New Relic 遥测数据并发送到 New Relic。一个 Connector 实例只处理一种遥测类型;同一插件提供三个可部署的具体类。com.newrelic.telemetry.events.EventsSinkConnector 把每条记录转换为一个 Event。记录值必须包含字符串字段 eventType,可选的 timestamp 使用 Unix 毫秒时间戳,其余受支持的字段作为事件属性。
com.newrelic.telemetry.logs.LogsSinkConnector 把每条记录转换为一条 Log。建议提供字符串字段 message;可选的 timestamp 使用 Unix 毫秒时间戳,其余受支持的字段作为日志属性。
com.newrelic.telemetry.metrics.MetricsSinkConnector 把每条记录转换为一个 Gauge、Count 或 Summary。记录值必须包含 name 和大小写敏感的 type;gauge、count 使用数值字段 value,summary 使用顶层字段 aggregated_summary.count、aggregated_summary.sum、aggregated_summary.min 和 aggregated_summary.max,dimensions 用于提供字符串维度。
使用无 Schema 的 JSON 时,记录值应转换为 Map;使用带 Schema 的数据时,记录值应转换为 Struct。带 Schema 的 timestamp 必须是 INT64,Gauge 和 Count 的 value 必须是 FLOAT64;Summary 的 aggregated_summary.count 必须是 INT32,其余三个聚合字段必须是 FLOAT64。复杂字段不会被递归展开;无 Schema Event 中的布尔值、对象、数组和 null 等值不会写入属性。Connector 会为遥测数据附加 Kafka Topic、Partition 和 Offset 元数据,但不会把 Kafka Key、Header 或记录时间戳自动映射为业务字段。
前置条件
- 准备可接收遥测数据的 New Relic 账号和 API Key,并确认账号所属区域是
US还是EU。
授权许可
使用 Apache License 2.0。快速开始
提前准备 Connect Cluster、Kafka、包含 Event 记录的 Topic 和 New Relic 账号,并确认网络连通和访问权限。具体准备和管理操作请参阅 管理 Connector。<new-relic-api-key> 替换为 New Relic API Key。如果账号位于欧盟区域,将 nr.region 改为 EU。示例显式使用无 Schema JSON,因此输入 telemetry-events Topic 的消息值可采用以下结构:
eventType 必须是非空字符串。省略 timestamp 时,Connector 在转换记录时使用当前时间。
配置
实例与遥测类型
name
Connector 实例的唯一名称。
- 类型:
string - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:使用不含控制字符的非空名称,并确保在 Connect Cluster 内唯一。
- 必填:是
connector.class
要加载的 New Relic 遥测 Connector 实现类。
- 类型:
string - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:Event 使用
com.newrelic.telemetry.events.EventsSinkConnector;Log 使用com.newrelic.telemetry.logs.LogsSinkConnector;Metric 使用com.newrelic.telemetry.metrics.MetricsSinkConnector。一个实例只能选择其中一个具体类。 - 必填:是
Kafka 输入
topics
要消费的 Kafka Topic 列表。
- 类型:
list - 默认值:
[] - 重要级别:高
- 有效值 / 注意事项:使用逗号分隔多个 Topic。必须在
topics和topics.regex中恰好配置一个。 - 必填:条件必填
topics.regex
用于动态匹配输入 Topic 的 Java 正则表达式。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:必须是有效的 Java 正则表达式。必须在
topics和topics.regex中恰好配置一个。 - 必填:条件必填
任务并行度
tasks.max
允许启动的最大 Task 数量。
- 类型:
int - 默认值:
1 - 重要级别:高
- 有效值 / 注意事项:最小值为
1。实际并行度受输入 Topic Partition 数量限制,每个 Task 使用独立的发送队列和 New Relic 客户端。
身份认证与区域
api.key
发送遥测数据时使用的 New Relic API Key。
- 类型:
password - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:必须使用有效的非空凭证。将其作为敏感信息管理,不要写入日志或提交到版本库。
- 必填:是
nr.region
New Relic 账号所在的数据区域。
- 类型:
string - 默认值:
US - 重要级别:低
- 有效值 / 注意事项:仅支持大小写敏感的
US或EU。建议始终显式配置,以确保 Task 使用预期区域。
网络连接
nr.client.timeout
New Relic HTTP 请求的调用超时时间,单位为毫秒。
- 类型:
int - 默认值:
2000 - 重要级别:低
- 有效值 / 注意事项:使用正整数;Connector 不校验取值范围。
nr.client.proxy.host
HTTP 代理主机名或地址。
- 类型:
string - 默认值:未设置(
null) - 重要级别:低
- 有效值 / 注意事项:只有同时配置
nr.client.proxy.host和nr.client.proxy.port时才启用代理。
nr.client.proxy.port
HTTP 代理端口。
- 类型:
int - 默认值:未设置(
null) - 重要级别:低
- 有效值 / 注意事项:只有同时配置
nr.client.proxy.host和nr.client.proxy.port时才启用代理。使用1到65535范围内的有效 TCP 端口。
批处理
nr.flush.max.records
单个遥测批次最多包含的记录数。
- 类型:
int - 默认值:
1000 - 重要级别:低
- 有效值 / 注意事项:使用正整数。达到该数量后,Connector 立即形成批次;该限制按记录数而不是请求字节数计算。
nr.flush.max.interval.ms
未达到最大记录数时,形成非空遥测批次的最长等待时间,单位为毫秒。
- 类型:
int - 默认值:
5000 - 重要级别:低
- 有效值 / 注意事项:使用正整数。较小值通常降低低流量场景的等待时间,但会增加请求频率。
数据转换与处理
key.converter
把 Kafka 消息 Key 反序列化为 Connect 数据的 Converter 类。
- 类型:
class - 默认值:未设置(
null) - 重要级别:低
- 有效值 / 注意事项:未设置时继承 Worker 级配置。New Relic Connector 不把 Key 映射为遥测业务字段。
value.converter
把 Kafka 消息值反序列化为 Connector 可处理的 Connect 数据。
- 类型:
class - 默认值:未设置(
null) - 重要级别:低
- 有效值 / 注意事项:未设置时继承 Worker 级配置。转换结果必须是带 Schema 的 Struct 或无 Schema 的 Map。无 Schema JSON 可使用
org.apache.kafka.connect.json.JsonConverter并设置value.converter.schemas.enable=false。
transforms
按顺序应用于记录的单消息转换 SMT 别名列表。
- 类型:
list - 默认值:
[] - 重要级别:低
- 有效值 / 注意事项:别名必须唯一,并通过
transforms.<alias>.type及该 SMT 的其他配置定义转换。可在记录进入 New Relic 类型转换前补充或整理所需字段。
最佳实践
首次接入时按遥测类型拆分实例
适用业务场景: 同一 Kafka 环境中同时存在事件、日志和指标流,需要让每类数据按自身字段结构进入 New Relic,并避免一种类型的异常记录影响其他类型。为每种遥测类型创建独立 Connector 实例,并订阅对应 Topic。 配置示例: 以下实例接收无 Schema JSON 日志:MetricsSinkConnector:
type 仅识别 gauge、count 和 summary;Summary 使用四个顶层点号字段 aggregated_summary.count、aggregated_summary.sum、aggregated_summary.min 和 aggregated_summary.max,不是嵌套对象。将 Topic、Converter 和消息生产端的数据契约作为一个整体管理,避免具体类与记录结构不匹配。
运行中按延迟要求调整批次
适用业务场景: Connector 已稳定发送数据,但低流量 Topic 的遥测可见时间偏长,或希望限制每次形成批次的记录数量。通过记录数和等待时间共同定义批次触发条件。 配置示例: 在快速开始配置基础上调整两个批处理阈值:500 条时可立即发送;未达到时,非空批次最多等待 2000 毫秒。示例值不是通用最优值,应结合输入速率、可接受延迟、请求频率和 Worker 内存持续观察并调整。Kafka Connect 的 Offset 提交周期不是 New Relic 批次触发条件。
积压增长时按 Partition 扩展 Task
适用业务场景: 输入 Topic 有多个 Partition,单个 Task 的消费能力不足且 New Relic 端仍有接收余量,需要增加并行消费来处理积压。 配置示例: 以下配置允许最多启动三个 Task:tasks.max。不同 Task 之间不提供全局记录顺序。
监控
监控内容
关注 Kafka Connect 健康状态、Connector 和 Task 状态、输入与发送吞吐、处理与目标可见延迟、消费组 Offset 提交和积压、转换或 HTTP 错误、SDK 重试以及 Worker JVM 堆内存、GC 和线程信号;仅在部署启用了相应 Kafka Connect 错误处理时关注 Converter 或 SMT 错误产生的 DLQ 活动。导入 Grafana 大盘
确认 Connect 指标已接入 Grafana 数据源,且采集标签满足大盘筛选条件;下载 Kafka Connect Dashboard,在 Grafana 中导入 JSON 并选择对应数据源。限制条件
- Events、Logs 和 Metrics 必须分别部署对应的具体 Connector 类;一个实例不能在运行时按记录内容切换遥测类型。
- Connector 异步发送遥测数据,Kafka Offset 可能在 New Relic 确认接收之前提交;它不提供端到端恰好一次保证,也不能无条件保证至少一次,故障或重启窗口中可能出现丢失或重复。
- Connector 内部的记录转换异常会被记录并跳过,不会进入 Kafka Connect DLQ;
errors.tolerance只影响记录进入 Connector 之前由 Converter 或 SMT 抛出的适用错误。 - 复杂对象和数组不会作为通用嵌套结构递归展开;生产端应按所选遥测类型提供受支持的属性、Metric
dimensions映射或 Summary 顶层点号字段。 - 异步重试和多 Task 并发可能改变记录到达 New Relic 的先后顺序,不应依赖跨批次或跨 Partition 的全局顺序。
常见问题
Connector 正常运行,但 New Relic 中没有出现数据
先确认connector.class 与消息值结构一致,api.key 和 nr.region 正确,并检查 Task 日志中的转换错误、HTTP 响应和重试信息。随后核对 Topic 是否持续有新记录、消费组 Offset 是否推进,以及批处理等待时间是否符合预期;修正字段类型或凭证后重新发送一条符合格式的新记录进行验证。
Metrics Connector 提示找不到类
Metric 的有效类名是com.newrelic.telemetry.metrics.MetricsSinkConnector。将 connector.class 修改为该完整类名后重新提交配置,并确认 Task 启动。
无 Schema JSON 记录出现类型转换错误
确认value.converter 使用 org.apache.kafka.connect.json.JsonConverter 且 value.converter.schemas.enable=false,并保证消息值是 JSON 对象而不是字符串化 JSON。Event 的 eventType、Log 的 message 应为字符串;Metric 的 name 和 type 必须存在,Gauge 或 Count 的 value 必须可解析为数值。
错误记录为什么没有进入 DLQ
New Relic 类型转换发生在 Connector 的put 处理内部,该路径会记录并跳过异常记录,而不会交给 Kafka Connect 的 DLQ 处理器。检查 Task 日志定位具体字段和类型错误,并在生产端、Converter 或 SMT 中修正记录;DLQ 仅可用于 Kafka Connect 在调用 Connector 前捕获的适用 Converter 或 SMT 错误。
安装的是 2.3.3,为什么运行时显示 2.3.0
该发布版本的 Connector 和 Task 运行时版本字符串仍返回2.3.0,因此插件信息或 New Relic 的 collector.version 属性可能显示该值。以实际安装的插件制品版本确认部署版本,不要仅根据运行时字符串降级或重复安装。