Skip to main content

概述

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 和大小写敏感的 typegaugecount 使用数值字段 valuesummary 使用顶层字段 aggregated_summary.countaggregated_summary.sumaggregated_summary.minaggregated_summary.maxdimensions 用于提供字符串维度。 使用无 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。必须在 topicstopics.regex 中恰好配置一个。
  • 必填:条件必填

topics.regex

用于动态匹配输入 Topic 的 Java 正则表达式。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:必须是有效的 Java 正则表达式。必须在 topicstopics.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
  • 重要级别:低
  • 有效值 / 注意事项:仅支持大小写敏感的 USEU。建议始终显式配置,以确保 Task 使用预期区域。

网络连接

nr.client.timeout

New Relic HTTP 请求的调用超时时间,单位为毫秒。
  • 类型int
  • 默认值2000
  • 重要级别:低
  • 有效值 / 注意事项:使用正整数;Connector 不校验取值范围。

nr.client.proxy.host

HTTP 代理主机名或地址。
  • 类型string
  • 默认值:未设置(null
  • 重要级别:低
  • 有效值 / 注意事项:只有同时配置 nr.client.proxy.hostnr.client.proxy.port 时才启用代理。

nr.client.proxy.port

HTTP 代理端口。
  • 类型int
  • 默认值:未设置(null
  • 重要级别:低
  • 有效值 / 注意事项:只有同时配置 nr.client.proxy.hostnr.client.proxy.port 时才启用代理。使用 165535 范围内的有效 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
Gauge 或 Count 消息值可采用以下结构:
关键说明: 不要在一个实例中混合三种记录结构。Metric 的 type 仅识别 gaugecountsummary;Summary 使用四个顶层点号字段 aggregated_summary.countaggregated_summary.sumaggregated_summary.minaggregated_summary.max,不是嵌套对象。将 Topic、Converter 和消息生产端的数据契约作为一个整体管理,避免具体类与记录结构不匹配。

运行中按延迟要求调整批次

适用业务场景: Connector 已稳定发送数据,但低流量 Topic 的遥测可见时间偏长,或希望限制每次形成批次的记录数量。通过记录数和等待时间共同定义批次触发条件。 配置示例: 在快速开始配置基础上调整两个批处理阈值:
关键说明: 批次达到 500 条时可立即发送;未达到时,非空批次最多等待 2000 毫秒。示例值不是通用最优值,应结合输入速率、可接受延迟、请求频率和 Worker 内存持续观察并调整。Kafka Connect 的 Offset 提交周期不是 New Relic 批次触发条件。

积压增长时按 Partition 扩展 Task

适用业务场景: 输入 Topic 有多个 Partition,单个 Task 的消费能力不足且 New Relic 端仍有接收余量,需要增加并行消费来处理积压。 配置示例: 以下配置允许最多启动三个 Task:
关键说明: 有效并行度不会超过可分配的 Topic Partition 数量。增加 Task 会创建独立的发送队列和客户端,并可能增加对 New Relic 的并发请求;先确认 Partition 数量和目标端容量,再逐步提高 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.keynr.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.JsonConvertervalue.converter.schemas.enable=false,并保证消息值是 JSON 对象而不是字符串化 JSON。Event 的 eventType、Log 的 message 应为字符串;Metric 的 nametype 必须存在,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 属性可能显示该值。以实际安装的插件制品版本确认部署版本,不要仅根据运行时字符串降级或重复安装。