概述
Adobe Experience Platform Sink Connector 消费 Kafka Topic 中的 JSON 消息,经 HTTP 流式采集入口批量发送到 Adobe Experience Platform(AEP)。它位于事件生产端与 AEP 数据流之间,适用于将客户活动、应用事件等数据接入已配置的 XDM Schema 和 Dataset,供后续分析或客户画像处理使用。 每条 Kafka 消息值对应一个 AEP 消息,Connector 将这些消息组成messages 数组发送。目标资源和字段映射取决于 AEP 数据流配置及消息内容。Connector 不自动创建 Dataset、生成映射或补充 AEP 消息头,也不会把 Kafka 消息键、分区、Offset 和 Headers 自动写入业务数据。
前置条件
- 在有权限的 AEP 项目、组织和 Sandbox 中准备目标 XDM Schema、Dataset,以及 HTTP API Source 的基础连接、源连接、目标连接和 Dataflow;取得基础连接返回的
inletUrl及对应的flowId。RAW 数据需要转换时,应在 AEP 配置 Data Prep Mapping 并关联到 Dataflow;已符合目标 Schema 的 XDM 数据可使用不带转换的路径。 - 确认 AEP 基础连接的
authenticationRequired策略。需要认证的入口必须具备有效且获授权的 Adobe 身份集成;涉及个人身份信息的数据应使用认证入口。资源管理 API 的 Bearer Token、API Key、IMS 组织和 Sandbox 请求头用于 AEP 资源管理,不应直接作为一套 Connector 配置照搬。 - 上游消息应具有与所选入口和数据流匹配的 AEP 消息结构及字段类型。创建 Dataflow 后等待约五分钟再发送数据,并考虑创建或更新数据流时的短暂采集暂停。
授权许可
使用 Apache License 2.0。快速开始
提前准备 Connect Cluster、Kafka 和上述 AEP 资源,确认网络连通和访问权限;创建和管理操作参见管理 Connector。以下最小配置适用于已预创建、authenticationRequired=false 的入口,以及消息键为 UTF-8 字符串或空值、消息值为普通 JSON 对象的 Topic;不适用于必须认证的入口或包含个人身份信息的数据。
aep.endpoint 使用 AEP 返回的 HTTPS inletUrl,路径应为 /collection/<connection-id>,不要填已经包含 /collection/batch/ 的地址;Connector 会自行转换为批量入口。无需显式填写默认的单 Task、批量阈值或关闭认证设置。value.converter.schemas.enable=false 使输入无需 Connect 的 schema / payload 外壳,但不能省略 AEP 所需的消息结构。
下面是上游生产者写入 Kafka 的单条消息值,展示 HTTP API Source 的 XDM 消息结构,不是 Connector 配置,也不是最终 HTTP 请求的外层 messages 数组。替换 Schema、Dataflow、Dataset 和业务字段后使用,xdmEntity 必须符合实际目标 Schema。
header 是 JSON 消息头,不是 HTTP Headers;flowId、datasetId、schemaRef 和 body 均不属于 Connector 参数。Connector 不会补齐这些字段,不要提前在每条 Kafka 消息外再套一层 messages。其他 AEP 采集路径可能要求不同的消息元数据,例如批量消息中的 imsOrgId;应以实际入口要求为准,不要将不同路径的样例混用。
应用配置后,先发送少量符合 Schema 的数据,分别检查 Connect 运行状态、采集响应,以及 AEP 验证和 Dataset 入库结果。HTTP 成功响应仅表示采集层的接受情况,不能证明最终 Dataset 入库成功。
安全警示:Connector 和 Task 启动时会在 INFO 日志中输出原始配置,包括可能存在的 Client Secret、授权码、代理密码或自定义 Authorization 请求头。仅关闭 DEBUG 或保持 errors.log.include.messages=false 不能解决此风险,也不能假定 Config Provider 替换后自动脱敏。在传入生产凭证前,应先落实启动日志脱敏或抑制方案及日志访问控制;敏感数据不要启用 DEBUG,凭证和消息内容不要附入工单或共享日志。
配置
实例与输入订阅
name
Connector 实例名称。
- 类型:
string - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:必填;非空且不含 ISO 控制字符,名称应唯一。
connector.class
选择 AEP Sink Connector。
- 类型:
string - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:必填;使用
com.adobe.platform.streaming.sink.impl.AEPSinkConnector。
topics
指定消费的 Kafka Topic。
- 类型:
list - 默认值:空列表
- 重要级别:高
- 有效值 / 注意事项:逗号分隔的 Topic 名称;与
topics.regex必须且只能选择一个非空配置,不得包含 DLQ Topic。
topics.regex
通过正则表达式订阅 Kafka Topic。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:有效的 Java 正则表达式;与非空
topics互斥,不得匹配 DLQ Topic。
tasks.max
允许创建的最大 Task 数量。
- 类型:
int - 默认值:
1 - 重要级别:高
- 有效值 / 注意事项:至少为
1。实际并行度受输入分区分配限制;各 Task 使用同一入口,不提供跨 Task 全局顺序。
消息转换
key.converter
反序列化 Kafka 消息键。
- 类型:
class - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:省略时使用 Worker Converter。须与上游键编码一致;快速开始针对 UTF-8 字符串键使用
org.apache.kafka.connect.storage.StringConverter,消息键不发送给 AEP。
value.converter
反序列化 Kafka 消息值。
- 类型:
class - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:省略时使用 Worker Converter。快速开始显式使用
org.apache.kafka.connect.json.JsonConverter,这不是默认类;输入必须符合 AEP 消息要求,而不只是语法合法的 JSON。
value.converter.schemas.enable
控制 JsonConverter 是否使用 Connect Schema 外壳。
- 类型:
boolean - 默认值:
true - 重要级别:高
- 有效值 / 注意事项:仅针对 JsonConverter。普通 JSON 对象使用
false;true要求 Connect 的schema/payload格式。此项不控制 AEP 的 XDM Schema 校验。
目标入口
aep.endpoint
指定 AEP 流式采集入口。
- 类型:
string(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:必填,无运行时回退;不能为空。填写完整 HTTPS 入口 URL,使用原始
/collection/路径;该路径会被替换为/collection/batch/,已经转换过的地址会被再次改写。
aep.connection.endpoint.headers
为采集请求附加 HTTP Headers。
- 类型:
string(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:值为请求头名称到字符串值的 JSON 对象文本。缺失或空字符串时运行时使用空映射;JSON 解析失败时不会添加这些请求头。此项不同于消息值中的
header,包含授权信息时有原始配置日志泄密风险。
批量发送
aep.flush.interval.seconds
控制一次消费批次内按经过时间触发发送的阈值,单位秒。
- 类型:
int(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:缺失或空白时运行时使用
1,按 Javaint运算乘以1000;乘法结果小于1时回退为1000ms。非法整数文本或超出int范围的文本会抛出NumberFormatException,不会回退。没有 ConfigDef 范围校验;建议使用正整数并避免乘法溢出。它不是后台定时器,消费批次结束时剩余消息仍会发送,不会跨消费批次等待凑满。
aep.flush.bytes.kb
控制一次消费批次内按累计消息大小触发发送的阈值,单位 KB。
- 类型:
int(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:缺失或空白时运行时使用
4,按 Javaint运算乘以1024;乘法结果小于1时回退为4096。非法整数文本或超出int范围的文本会抛出NumberFormatException,不会回退。没有 ConfigDef 范围校验;建议使用正整数并避免乘法溢出。实际累计的是 Java 字符串长度,不是 UTF-8 字节数,也不包含最终请求包装开销。加入记录后才检查阈值,不能以此确保请求满足 Adobe 的 1 MB 大小限制。
HTTP 连接与重试
aep.connection.timeout
采集 HTTP 请求的连接超时,单位毫秒。
- 类型:
int(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:缺失或空白时运行时使用
5000;使用正整数。只控制采集连接,不是 Token 请求超时设置。
aep.connection.readTimeout
采集 HTTP 请求的读取超时,单位毫秒。
- 类型:
int(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:缺失或空白时运行时使用
60000;使用正整数,保留名称大小写。不用于 Token 请求。
aep.connection.maxRetries
限制单次采集 HTTP 调用的总尝试次数。
- 类型:
int(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:缺失或空白时运行时使用
3,表示最多三次总尝试,不是初次请求再加三次重试。没有参数范围校验;要执行发送应使用正数,不要设置为零或负数。对 5xx 和 I/O 异常按固定间隔重新发送同一请求体;429 等其他非 2xx 不走该重试路径,响应内容解析失败也不在此范围。重发可能造成重复。
aep.connection.retryBackoff
采集 HTTP 尝试失败后的固定等待时长,单位毫秒。
- 类型:
int(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:缺失或空白时运行时使用
300;使用非负整数。最后一次 5xx 或 I/O 失败后也可能等待;等待会阻塞 Task,不是认证重试或自适应限速设置。
身份认证
aep.connection.auth.enabled
启用 Connector 的 Bearer Token 获取与附加功能。
- 类型:
string(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:缺失时运行时使用字符串
false。只有精确的小写true才启用;TRUE或带空白的值不会启用。此项不能改变 AEP 入口的认证策略。
aep.connection.auth.token.type
选择 Adobe Token 获取方式。
- 类型:
string(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:启用认证且未配置时运行时使用
jwt_token;支持精确值access_token、jwt_token、oauth2_access_token,空值或未知值会失败。JWT 认证路径已弃用;迁移时应按 Adobe 身份集成要求显式选择 OAuth2 并确认授权和 scope 兼容性,不能因选项存在就假定旧 JWT 仍获服务端支持。
aep.connection.auth.client.id
指定 Adobe 身份集成的 Client ID。
- 类型:
string(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:无运行时回退;三种认证方式启用时均必填。它不是 Access Token。
aep.connection.auth.client.secret
指定 Adobe 身份集成的 Client Secret。
- 类型:
string(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:无运行时回退;三种认证方式启用时均必填。此项没有 PASSWORD 类型的自动脱敏保护,传入前必须解决原始配置日志泄密风险。
aep.connection.auth.endpoint
指定 Adobe IMS 基础 URL。
- 类型:
string(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:无运行时回退;启用认证时必填,不能依靠环境变量省略。三种方式分别附加
/ims/token/v1、/ims/exchange/jwt/、/ims/token/v3,不要提前附加这些路径。
aep.connection.auth.client.code
指定 IMS 授权码交换所需的授权码。
- 类型:
string(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:无运行时回退;仅
access_token方式必填,其他方式不使用。它不是预签发的 Bearer Token;属于敏感凭证,可能出现在原始配置日志中。
旧版 JWT 认证
aep.connection.auth.imsOrg
指定 JWT 认证的 IMS 组织 ID。
- 类型:
string(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:无运行时回退;仅旧版
jwt_token方式必填,其他认证方式不使用。该认证路径已弃用,替代路径为符合 Adobe 身份集成要求的oauth2_access_token。
aep.connection.auth.accountKey
指定 JWT 认证的技术账号标识。
- 类型:
string(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:无运行时回退;仅旧版
jwt_token方式必填,不是 RSA 私钥。该认证路径已弃用;其他方式不使用,迁移时改用对应的 OAuth2 Client 配置。
aep.connection.auth.filePath
指定 JWT 认证使用的 RSA 私钥文件路径。
- 类型:
string(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:无运行时回退;仅旧版
jwt_token方式必填。每个承担 Task 的 Worker 都必须能读取该本地文件;使用 PEM 编码的 PKCS#8 私钥,文件不得超过 1 MiB。不要将私钥正文写入配置。该认证路径已弃用,OAuth2 方式不使用此项。
HTTP 代理
aep.connection.proxy.host
指定代理主机,供采集请求及认证提供方使用。
- 类型:
string(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:缺失或空白时运行时使用
null,不指定代理主机;代理行为由 HTTP 连接处理。
aep.connection.proxy.port
指定代理端口。
- 类型:
int(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:缺失或空白时运行时使用
443。配置代理主机时填写有效的 TCP 端口;不存在 Connector 级范围校验。
aep.connection.proxy.user
指定需要认证的代理用户名。
- 类型:
string(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:缺失或空白时运行时使用
null;按代理要求与密码配合。代理认证使用 JVM 全局默认 Authenticator,需考虑同 Worker 内其他连接的影响。
aep.connection.proxy.password
指定代理认证密码。
- 类型:
string(运行时读取类型) - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:缺失或空白时运行时使用
null。没有自动 PASSWORD 脱敏保护;启动原始配置日志和代理 DEBUG 日志均存在凭证暴露风险。
错误处理
errors.tolerance
控制框架可处理记录错误的容忍策略。
- 类型:
string - 默认值:
none - 重要级别:中
- 有效值 / 注意事项:
none或all。all不代表所有远端 HTTP 错误均可恢复,也不能防止 Connector 内部未报告的记录丢弃;401/403 仍可能导致 Task 失败。
errors.deadletterqueue.topic.name
指定 Sink 错误记录的 DLQ Topic。
- 类型:
string - 默认值:空字符串
- 重要级别:中
- 有效值 / 注意事项:空值禁用 DLQ;非空时必须从输入订阅排除。只有进入框架错误报告路径的记录才可能写入,不能将其视为全部 AEP 失败记录的备份。
errors.deadletterqueue.topic.replication.factor
设置自动创建 DLQ Topic 时的副本因子。
- 类型:
short - 默认值:
3 - 重要级别:中
- 有效值 / 注意事项:针对尚不存在的 DLQ Topic,应与 Kafka 部署可支持的副本数匹配;不改变已有 Topic 的副本数。
errors.deadletterqueue.context.headers.enable
为 DLQ 记录附加错误上下文请求头。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:启用时添加
__connect.errors.前缀的 Headers;需评估上下文信息暴露风险。这些 Headers 不属于发往 AEP 的消息头。
errors.log.enable
控制框架错误日志。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:独立于 Connector 日志级别;关闭时也不能阻止 Connector 启动日志输出原始配置。
errors.log.include.messages
控制框架错误日志是否包含消息详情。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:敏感输入应保持关闭;此项不是 Adobe 配置或代理凭证的脱敏机制。
最佳实践
积压增长时增加消费并行度
适用业务场景:已完成接入并确认数据能正常写入 Dataset,但持续运行时单 Task 的消费积压不断增长。希望利用输入 Topic 的多个分区并行发送,同时保持消息格式和目标入口不变。 配置示例:在快速开始的配置上添加以下设置。这里的2 是扩展起点示例,不是通用最优值;输入 Topic 应具有足够的可分配分区。
监控
监控内容
关注 Kafka Connect 集群健康、Connector / Task 状态、消费吞吐与积压、端到端延迟、Offset 提交、错误和重试,以及 Worker JVM 的内存、GC 和线程信号;只有配置了相关错误处理时才关注 DLQ 活动。RUNNING、HTTP 2xx 或已提交 Offset 均不代表 Dataset 入库成功,还需独立观察 AEP Dataflow、Dataset 批次状态和流式验证结果,将采集层接受量与最终入库结果分别核对。导入 Grafana 大盘
下载共享的 Kafka Connect Grafana Dashboard,确认已采集 Connect / Worker 指标并配置兼容的数据源及大盘所需的集群、Connector、Task 等标签,然后在 Grafana 导入 JSON 并选择对应数据源。限制条件
- Connector 不自动创建 AEP 资源、执行 Data Prep 映射或补充 XDM 消息结构;也不提供按 Topic 配置不同目标入口的路由。
- 不提供 CDC 删除、更新合并或 Schema 演进管理;Kafka 空值和 tombstone 不代表 AEP 删除操作,可能触发数据转换异常。
- 原始 JSON 字符串解析失败时可能被直接排除且不进入错误报告路径;部分目标端错误也可能在记录日志后继续推进 Offset,不能将 Offset 推进解释为完整交付。
- 不提供事务写入、幂等键或自动去重;HTTP 重试、Offset 提交前的重启及重平衡可能产生重复,不保证无条件至少一次送达 AEP,也不保证恰好一次。
- 批量阈值不等于最终 HTTP 请求字节大小,不能自动确保 Adobe 的 1 MB 请求上限;单条大消息也可能超过阈值。
- 不提供全局事件顺序或 AEP 最终入库顺序保障,各 Task 和分区可能并行处理。
常见问题
Task 显示 RUNNING,为什么 Dataset 中仍没有数据?
启动状态不是入口连接检查,HTTP 接受也不等于 Dataset 入库。确认入口地址、Dataflow 状态、消息内的flowId 和 datasetId 是否对应,再检查目标 Schema、XDM 字段类型和 AEP 流式验证及 Dataset 批次错误。新建 Dataflow 后等待约五分钟再发送。不要只凭 Offset 推进判断成功,应使用可识别的业务事件核对最终入库结果。
为什么收到 401 或 403 后没有自动恢复?
检查 AEP 入口是否要求认证,以及 Token 是否有效、身份集成是否拥有所需权限。启用 Connector 认证必须使用精确的小写true,填写选定方式必需的参数;不要依赖隐式选择旧版 JWT。OAuth2 按 Token 过期时间缓存和刷新,但不会因 401 自动清除缓存并重新获取 Token 后重发。处理授权问题后再恢复 Task;排查时先确保日志不会暴露凭证。
HTTP 返回成功,为什么仍可能有部分消息失败?
批量请求的 207 响应可以包含逐条结果,不能只检查 HTTP 状态码。Connector 在status 字段存在且值非 null 时判断为失败,不要求字符串非空,并从 xactionId 最后一个连字符后的内容获取记录索引;Adobe 公开样例则使用 statusCode,成功标识可能以冒号分隔,失败项还可能没有 xactionId。这是样例与插件解析规则的静态契约差异,并不代表已经确认真实云端发生故障。上线前应核对实际入口响应,尤其不能依赖插件识别只含 statusCode 的失败项;它也不会自动重试失败的单条消息。结合 AEP 验证和 Dataset 状态核对结果,必要时联系维护方确认响应兼容性,获取响应时应避免泄露业务数据。
配置了 DLQ,为什么没有找到所有失败消息?
DLQ 只覆盖被交给框架错误报告路径的记录,不覆盖所有内部解析丢弃或所有远端响应问题;启用错误容忍也不会解决 AEP 入库后的验证失败。检查错误容忍策略及输入订阅是否排除了 DLQ,再确认 DLQ Topic 已存在,或具备自动创建所需的权限和副本数条件。结合 Task 日志和 AEP 结果核对,不要把没有 DLQ 记录当作全部消息成功,也不要把记录报告等同于同步持久化完成。为什么恢复或重启后出现重复数据?
请求可能已被 AEP 接受,但连接中断或 Kafka Offset 尚未提交,随后 HTTP 重试或恢复消费会重新发送。Connector 不自动去重,应按业务事件标识核对重复,并在目标数据处理链路单独设计、确认去重规则;不能假定 XDM 的_id 字段本身就让该 Connector 获得恰好一次语义。