CloudEvents
导出 CloudEvents
CloudEvents 是一种以通用方式描述事件数据的规范,其目标是在服务、平台和系统之间实现互操作性。Debezium 允许你配置 Db2、Informix、MongoDB、MySQL、Oracle、PostgreSQL 或 SQL Server 连接器,以发出符合 CloudEvents 规范的变更事件记录。
| CloudEvents 的支持目前处于孵化状态。这意味着确切的语义、配置选项和其他细节可能会在未来的修订版本中根据反馈发生变化。如果你有特定的需求,或者在使用此功能时遇到任何问题,请告诉我们。 |
|---|
CloudEvents 规范定义了:
- 一组标准化的事件属性
- 定义自定义属性的规则
- 将事件格式映射到 JSON 或 Apache Avro 等序列化表示形式的编码规则
- 面向 Apache Kafka、HTTP 或 AMQP 等传输层的协议绑定
要配置 Debezium 连接器以发出符合 CloudEvents 规范的变更事件记录,Debezium 提供了 io.debezium.converters.CloudEventsConverter,它是一个 Kafka Connect 消息转换器。
目前,仅支持结构化映射模式。CloudEvents 变更事件封装(envelope)格式可以是 JSON 或 Avro,并且对于每种封装格式,你都可以使用 JSON 或 Avro 作为 data 格式。预计未来的 Debezium 版本将支持二进制映射模式。
有关使用 Avro 的信息,请参阅:
事件格式示例
以下示例展示了 PostgreSQL 连接器发出的 CloudEvents 变更事件记录的样子。在此示例中,PostgreSQL 连接器被配置为使用 JSON 作为 CloudEvents 格式的封装,同时也作为 data 格式。
{
"id" : "name:test_server;lsn:29274832;txId:565",
"source" : "/debezium/postgresql/test_server",
"specversion" : "1.0",
"type" : "io.debezium.connector.postgresql.DataChangeEvent",
"time" : "2020-01-13T13:55:39.738Z",
"datacontenttype" : "application/json",
"iodebeziumop" : "r",
"iodebeziumversion" : "3.6.3.Final",
"iodebeziumconnector" : "postgresql",
"iodebeziumname" : "test_server",
"iodebeziumtsms" : "1578923739738",
"iodebeziumsnapshot" : "true",
"iodebeziumdb" : "postgres",
"iodebeziumschema" : "s1",
"iodebeziumtable" : "a",
"iodebeziumlsn" : "29274832",
"iodebeziumxmin" : null,
"iodebeziumtxid": "565",
"iodebeziumtxtotalorder": "1",
"iodebeziumtxdatacollectionorder": "1",
"data" : {
"before" : null,
"after" : {
"pk" : 1,
"name" : "Bob"
}
}
}以下列表描述了前述 CloudEvents 变更事件记录中的部分字段:
id
连接器根据变更事件的内容为该变更事件生成的唯一 ID。
source
事件的来源,即连接器配置中 topic.prefix 属性所指定的数据库逻辑名称。
specversion
CloudEvents 规范版本。
type
生成该变更事件的连接器类型。此字段的格式为 io.debezium.connector.CONNECTOR_TYPE.DataChangeEvent。
CONNECTOR_TYPE 的有效取值包括 db2、mariadb,、informix、mongodb、mysql、oracle、postgresql 或 sqlserver。
time
源数据库中发生变更的时间。
datacontenttype
描述 data 属性的内容类型。可能的取值有 json(如本示例所示)或 avro。
iodebeziumop
操作标识符。可能的取值有 r(读取)、c(创建)、u(更新)或 d(删除)。
iodebeziumversion
Debezium 变更事件中已知的所有 source 属性,均通过在属性名称前加上 iodebezium 前缀映射为 CloudEvents 扩展属性。
iodebeziumtxid
在连接器中启用该功能后,Debezium 变更事件中已知的每个 transaction 属性,均通过在属性名称前加上 iodebeziumtx 前缀映射为 CloudEvents 扩展属性。
data
实际的数据变更。根据操作类型和连接器的不同,数据可能包含 before、after 或 patch 字段。
下面的示例展示了 PostgreSQL 连接器发出的 CloudEvents 变更事件记录的样子。在此示例中,PostgreSQL 连接器再次被配置为使用 JSON 作为 CloudEvents 格式信封,但这一次连接器被配置为对 data 格式使用 Avro。
{
"id" : "name:test_server;lsn:33227720;txId:578",
"source" : "/debezium/postgresql/test_server",
"specversion" : "1.0",
"type" : "io.debezium.connector.postgresql.DataChangeEvent",
"time" : "2020-01-13T14:04:18.597Z",
"datacontenttype" : "application/avro",
"dataschema" : "http://my-registry/schemas/ids/1",
"iodebeziumop" : "r",
"iodebeziumversion" : "3.6.3.Final",
"iodebeziumconnector" : "postgresql",
"iodebeziumname" : "test_server",
"iodebeziumtsms" : "1578924258597",
"iodebeziumsnapshot" : "true",
"iodebeziumdb" : "postgres",
"iodebeziumschema" : "s1",
"iodebeziumtable" : "a",
"iodebeziumtxId" : "578",
"iodebeziumlsn" : "33227720",
"iodebeziumxmin" : null,
"iodebeziumtxid": "578",
"iodebeziumtxtotalorder": "1",
"iodebeziumtxdatacollectionorder": "1",
"data" : "AAAAAAEAAgICAg=="
}以下列表描述了使用 Avro 格式的 CloudEvents 变更事件记录中的部分字段:
datacontenttype
表示 data 属性包含 Avro 二进制数据。
dataschema
Avro 数据所遵循的模式的 URI。
data
data 属性包含经过 base64 编码的 Avro 二进制数据。
也可以对信封和 data 属性都使用 Avro。
示例配置
在 Debezium 连接器配置中配置 io.debezium.converters.CloudEventsConverter。以下示例展示了如何配置 CloudEvents 转换器,以输出具备以下特征的变更事件记录:
- 使用 JSON 作为信封。
- 使用位于
http://my-registry/schemas/ids/1的模式注册表,将data属性序列化为二进制 Avro 数据。
...
"value.converter": "io.debezium.converters.CloudEventsConverter",
"value.converter.serializer.type" : "json",
"value.converter.data.serializer.type" : "avro",
"value.converter.avro.schema.registry.url": "http://my-registry/schemas/ids/1"
...以下列表说明了前述 CloudEvents 转换器配置中的部分字段:
value.converter.serializer.type
指定 serializer.type 是可选的,因为 json 是默认值。
CloudEvents 转换器转换 Kafka 记录的值。在同一个连接器配置中,如果你要处理记录的键,可以指定 key.converter。例如,你可以指定 StringConverter、LongConverter、JsonConverter 或 AvroConverter。
元数据来源及部分 CloudEvents 字段的配置
默认情况下,metadata.source 属性由五部分组成,如下例所示:
"value,id:generate,type:generate,traceparent:header,dataSchemaName:generate"第一部分指定获取记录元数据的来源;允许的取值为 value 和 header。接下来的各部分指定转换器如何为以下 CloudEvents 字段和数据模式名称填充值:
idtypetraceparentdataSchemaName(该模式在 Schema Registry 中注册时所使用的名称)
转换器可以使用以下两种方式之一来填充每个字段:
generate
转换器为该字段生成一个值。
header
转换器从消息头部中获取该字段的值。
traceparent CloudEvents 字段的值只能从消息头部中获取。 |
|---|
获取记录元数据
要构造 CloudEvent,转换器需要来源、操作和事务元数据。通常,转换器可以从记录的值中检索这些元数据。但在某些情况下,在转换器接收到记录之前,记录可能已被处理,导致其值中不再包含元数据,例如,记录经过 Outbox Event Router SMT 处理之后就是这种情况。为了保留所需的元数据,你可以采用以下方式,通过记录头部来传递元数据。
操作步骤
- 实现一种机制,在记录到达转换器之前将其元数据记录到记录头部中,例如使用
HeaderFromSMT。 - 将转换器的
metadata.source属性值设置为header。
以下示例展示了使用 Outbox Event Router SMT 和 HeaderFrom SMT 的连接器配置:
...
"tombstones.on.delete": false,
"transforms": "addMetadataHeaders,outbox",
"transforms.addMetadataHeaders.type": "org.apache.kafka.connect.transforms.HeaderFrom$Value",
"transforms.addMetadataHeaders.fields": "source,op,transaction",
"transforms.addMetadataHeaders.headers": "source,op,transaction",
"transforms.addMetadataHeaders.operation": "copy",
"transforms.addMetadataHeaders.predicate": "isHeartbeat",
"transforms.addMetadataHeaders.negate": true,
"transforms.outbox.type": "io.debezium.transforms.outbox.EventRouter",
"transforms.outbox.table.expand.json.payload": true,
"transforms.outbox.table.fields.additional.placement": "type:header",
"predicates": "isHeartbeat",
"predicates.isHeartbeat.type": "org.apache.kafka.connect.transforms.predicates.TopicNameMatches",
"predicates.isHeartbeat.pattern": "__debezium-heartbeat.*",
"value.converter": "io.debezium.converters.CloudEventsConverter",
"value.converter.metadata.source": "header",
"header.converter": "org.apache.kafka.connect.json.JsonConverter",
"header.converter.schemas.enable": true
...要使用 HeaderFrom 转换,可能需要过滤墓碑消息和心跳消息。 |
|---|
metadata.source 属性的 header 值是一项全局设置。因此,即使你省略了某个属性值的部分内容(例如 id 和 type 来源),转换器也会为被省略的部分生成 header 值。
获取 CloudEvent 元数据
默认情况下,CloudEvents 转换器会自动生成 CloudEvent 的 id 和 type 字段的值,并为其 data 字段生成架构名称。只有在 opentelemetry.tracing.attributes.enable 设置为 true 时,消息中才会包含 traceparent CloudEvents 字段。你可以通过修改默认值并在相应的消息头中指定字段值,来自定义转换器填充这些字段的方式。例如:
"value.converter.metadata.source": "value,id:header,type:header,traceparent:header,dataSchemaName:header"在上述配置生效后,你可以配置上游函数,添加 id、type、traceparent 和 dataSchemaName 头部,并指定你希望传递给 CloudEvents 转换器的值。
如果你只想为 id 头部提供值,请使用:
"value.converter.metadata.source": "value,id:header,type:generate,traceparent:header,dataSchemaName:generate"要配置转换器从消息头中获取 id、type、traceparent 和 dataSchemaName 元数据,请使用以下简写语法:
"value.converter.metadata.source": "header"要让转换器能够从消息头中获取数据模式名称,必须将 schema.data.name.source.header.enable 设置为 true。
配置选项
配置 Debezium 连接器以使用 CloudEvent 转换器时,可以指定以下选项:
表 1. CloudEvents 转换器配置选项说明
选项
默认值
说明
json
用于 CloudEvents 信封结构的编码类型。取值可以为 json 或 avro。
json
用于 data 属性的编码类型。取值可以为 json 或 avro。
不适用
使用 JSON 时要传递给底层转换器的任何配置选项。json. 前缀会被移除。
不适用
使用 Avro 时要传递给底层转换器的任何配置选项。avro. 前缀会被移除。例如,对于 Avro 格式的 data,应指定 avro.schema.registry.url 选项。
none
指定应如何调整模式名称以兼容连接器所使用的消息转换器。取值可以为 none 或 avro。
none
指定模式在 Schema Registry 中注册时所使用的 CloudEvents 模式名称。当 serializer.type 为 json 时,该设置会被忽略,因为此时记录的值是无模式的。如果未指定此属性,则使用默认算法生成模式名称:${serverName}.${databaseName}.CloudEvents.Envelope。
schema.data.name.source.header.enable
false
指定转换器是否可以从消息头中获取 CloudEvents data 字段的模式名称。该模式名称取自 metadata.source 属性中指定的 dataSchemaName 参数。
opentelemetry.tracing.attributes.enable
false
指定转换器生成云事件时是否包含 OpenTelemetry 追踪属性。取值可以为 true 或 false。
true
指定转换器生成云事件时是否包含扩展属性。取值可以为 true 或 false。
value,id:generate,type:generate,traceparent:header,dataSchemaName:generate
一个以逗号分隔的列表,用于指定转换器从何处检索 id、type 和 traceparent CloudEvents 字段以及 dataSchemaName 参数的元数据值(源、操作、事务),其中 dataSchemaName 参数指定该模式(schema)在 Schema Registry 中注册时所使用的名称。列表中的第一个元素是全局设置,用于指定元数据的来源。元数据的来源可以是 value 或 header。全局设置之后是一系列键值对。每对中的第一个元素指定某个 CloudEvent 字段的名称(id、type 或 traceparent),或数据模式的名称(dataSchemaName)。每对中的第二个元素指定转换器填充该字段值的方式。有效值为 generate 或 header。每对中的两个值用冒号分隔,例如:
value,id:header,type:generate,traceparent:header,dataSchemaName:header
有关配置示例,请参阅元数据来源和部分 CloudEvents 字段的配置。
评论
登录后参与评论
KnowForge