-
Notifications
You must be signed in to change notification settings - Fork 1.2k
cdc: add kafka max-message-bytes related explain #21832
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
5f4ddf3
ce6fec4
c4edd7d
c9a7de2
a882d2f
ecf273d
a7e87c1
7953c42
2361e1d
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change | ||||||||
|---|---|---|---|---|---|---|---|---|---|---|
|
|
@@ -44,6 +44,19 @@ Value: | |||||||||
| * 长度及协议版本号均为大端序 int64 类型 | ||||||||||
| * 当前协议版本号为 `1` | ||||||||||
|
|
||||||||||
| ### 控制 Message 中的 Event 数量和大小 | ||||||||||
|
|
||||||||||
| Open Protocol 可将一个或多个 Row Changed Event 编码为一条 Message。Kafka Sink 的以下参数分别控制一条 Message 中的 Row Changed Event 数量和 Message 大小: | ||||||||||
|
|
||||||||||
| | 参数 | 作用 | | ||||||||||
| | --- | --- | | ||||||||||
| | `max-batch-size` | 每条 Message 最多包含的 Row Changed Event 数量,默认值为 `16`。 | | ||||||||||
| | `max-message-bytes` | 控制 Message 的大小阈值。 | | ||||||||||
|
|
||||||||||
| Kafka Sink 会比较 changefeed 中配置的 `max-message-bytes` 与 Kafka 允许的消息大小限制,并使用其中的较小值作为实际阈值,避免批量编码后的 Message 超过 Kafka 允许的大小。Kafka 消息大小限制的确定方式参见[Kafka 消息大小限制](/ticdc/ticdc-sink-to-kafka.md#kafka-消息大小限制)。 | ||||||||||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Suggested change
|
||||||||||
|
|
||||||||||
| 当 Message 中的 Row Changed Event 数量已经达到 `max-batch-size`,或者继续加入下一个 Row Changed Event 会使 Message 超过实际生效的大小阈值时,后续 Row Changed Event 会被写入一条新的 Message。 | ||||||||||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Suggested change
|
||||||||||
|
|
||||||||||
| ## Event 格式定义 | ||||||||||
|
|
||||||||||
| 本部分介绍 Row Changed Event、DDL Event 和 Resolved Event 的格式定义。 | ||||||||||
|
|
||||||||||
| Original file line number | Diff line number | Diff line change | ||||
|---|---|---|---|---|---|---|
|
|
@@ -77,13 +77,13 @@ URI 中可配置的的参数如下: | |||||
| | `kafka-version` | 下游 Kafka 版本号。该值需要与下游 Kafka 的实际版本保持一致。 | | ||||||
| | `kafka-client-id` | 指定同步任务的 Kafka 客户端的 ID(可选,默认值为 `TiCDC_sarama_producer_同步任务的 ID`)。 | | ||||||
| | `partition-num` | 下游 Kafka partition 数量(可选,不能大于实际 partition 数量,否则创建同步任务会失败,默认值 `3`)。| | ||||||
| | `max-message-bytes` | 每次向 Kafka broker 发送消息的最大数据量(可选,默认值 `10MB`,最大值为 `100MB`)。从 v5.0.6 和 v4.0.6 开始,默认值分别从 `64MB` 和 `256MB` 调整至 `10MB`。| | ||||||
| | `max-message-bytes` | Open Protocol Message 的大小阈值(可选,默认值 `10 MB`,最大值为 `100 MB`)。实际生效规则参见[控制 Message 中的 Event 数量和大小](/ticdc/ticdc-open-protocol.md#控制-message-中的-event-数量和大小)。| | ||||||
|
3AceShowHand marked this conversation as resolved.
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Suggested change
|
||||||
| | `max-batch-size` | Open Protocol Message 最多包含的 Row Changed Event 数量(可选,默认值 `16`)。详细说明参见[控制 Message 中的 Event 数量和大小](/ticdc/ticdc-open-protocol.md#控制-message-中的-event-数量和大小)。| | ||||||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Suggested change
|
||||||
| | `replication-factor` | Kafka 消息保存副本数(可选,默认值 `1`),需要大于等于 Kafka 中 [`min.insync.replicas`](https://kafka.apache.org/33/documentation.html#brokerconfigs_min.insync.replicas) 的值。 | | ||||||
| | `required-acks` | 在 `Produce` 请求中使用的配置项,用于告知 broker 需要收到多少副本确认后才进行响应。可选值有:`0`(`NoResponse`:不发送任何响应,只有 TCP ACK),`1`(`WaitForLocal`:仅等待本地提交成功后再响应)和 `-1`(`WaitForAll`:等待所有同步副本提交后再响应。最小同步副本数量可通过 broker 的 [`min.insync.replicas`](https://kafka.apache.org/33/documentation.html#brokerconfigs_min.insync.replicas) 配置项进行配置)。(可选,默认值为 `-1`)。 | | ||||||
| | `compression` | 设置发送消息时使用的压缩算法(可选值为 `none`、`lz4`、`gzip`、`snappy` 和 `zstd`,默认值为 `none`)。注意 Snappy 压缩文件必须遵循[官方 Snappy 格式](https://github.com/google/snappy)。不支持其他非官方压缩格式。| | ||||||
| | `auto-create-topic` | 当传入的 `topic-name` 在 Kafka 集群不存在时,TiCDC 是否要自动创建该 topic(可选,默认值 `true`)。 | | ||||||
| | `enable-tidb-extension` | 可选,默认值是 `false`。当输出协议为 `canal-json` 时,如果该值为 `true`,TiCDC 会发送 [WATERMARK 事件](/ticdc/ticdc-canal-json.md#watermark-event),并在 Kafka 消息中添加 TiDB 扩展字段。从 6.1.0 开始,该参数也可以和输出协议 `avro` 一起使用。如果该值为 `true`,TiCDC 会在 Kafka 消息中添加[三个 TiDB 扩展字段](/ticdc/ticdc-avro-protocol.md#tidb-扩展字段)。| | ||||||
| | `max-batch-size` | 从 v4.0.9 开始引入。当消息协议支持把多条变更记录输出至一条 Kafka 消息时,该参数用于指定这一条 Kafka 消息中变更记录的最多数量。目前,仅当 Kafka 消息的 `protocol` 为 `open-protocol` 时有效(可选,默认值 `16`)。| | ||||||
| | `enable-tls` | 连接下游 Kafka 实例是否使用 TLS(可选,默认值 `false`)。 | | ||||||
| | `ca` | 连接下游 Kafka 实例所需的 CA 证书文件路径(可选)。 | | ||||||
| | `cert` | 连接下游 Kafka 实例所需的证书文件路径(可选)。 | | ||||||
|
|
@@ -108,16 +108,11 @@ URI 中可配置的的参数如下: | |||||
|
|
||||||
| ### 最佳实践 | ||||||
|
|
||||||
| * TiCDC 推荐用户自行创建 Kafka Topic,你至少需要设置该 Topic 每次向 Kafka broker 发送消息的最大数据量和下游 Kafka partition 的数量。在创建 changefeed 的时候,这两项设置分别对应 `max-message-bytes` 和 `partition-num` 参数。 | ||||||
| * TiCDC 推荐用户自行创建 Kafka Topic,并根据业务需要设置该 Topic 的 `max.message.bytes` 和 partition 数量。创建 changefeed 时,可以通过 `partition-num` 指定 partition 数量。 | ||||||
| * 如果你在创建 changefeed 时,使用了尚未存在的 Topic,那么 TiCDC 会尝试使用 `partition-num` 和 `replication-factor` 参数自行创建 Topic,建议明确指定这两个参数。 | ||||||
| * 在大多数情况下,建议使用 `canal-json` 协议。 | ||||||
| * 如果 TiCDC 上游的数据变更很少,比如可能会出现超过 10 分钟没有数据变更的情况,建议在 Kafka broker 的配置文件中调大 Kafka 的连接空闲超时时间,详情参考[为什么 TiCDC 同步到 Kafka 的任务经常因 `broken pipe` 报错而失败](/ticdc/ticdc-faq.md#为什么-ticdc-同步到-kafka-的任务经常因-broken-pipe-报错而失败)。 | ||||||
|
|
||||||
| > **注意:** | ||||||
| > | ||||||
| > 当 `protocol` 为 `open-protocol` 时,TiCDC 会将多个事件编码到同一个 Kafka 消息中,并尽量避免在此过程中生成长度超过 `max-message-bytes` 的消息。 | ||||||
| > 如果单条数据变更编码得到的消息大小超过了 `max-message-bytes` 个字节,changefeed 会报错,并打印错误日志。 | ||||||
|
|
||||||
| ### TiCDC 使用 Kafka 的认证与授权 | ||||||
|
|
||||||
| 使用 Kafka 的 SASL 认证时配置样例如下所示: | ||||||
|
|
@@ -162,13 +157,12 @@ URI 中可配置的的参数如下: | |||||
| - 对 Cluster 资源类型的 `DescribeConfig` 权限。 | ||||||
|
|
||||||
| 各权限的使用场景如下: | ||||||
|
|
||||||
| | 资源类型 | 操作类型 | 使用场景 | | ||||||
| | :-------------| :------------- | :--------------------------------| | ||||||
| | Cluster | `DescribeConfig` | Changefeed 运行过程中,获取集群元数据 | | ||||||
| | Topic | `Describe` | Changefeed 启动时,尝试创建 Topic | | ||||||
| | Topic | `Create` | Changefeed 启动时,尝试创建 Topic | | ||||||
| | Topic | `Write` | 发送数据到 Topic | | ||||||
| | :-------------| :--------------- | :--------------------------------| | ||||||
| | Cluster | `DescribeConfig` | Changefeed 运行过程中,获取集群元数据 | | ||||||
| | Topic | `Describe` | Changefeed 启动时,尝试创建 Topic | | ||||||
| | Topic | `Create` | Changefeed 启动时,尝试创建 Topic | | ||||||
| | Topic | `Write` | 发送数据到 Topic | | ||||||
|
|
||||||
| 创建或启动 Changefeed 时,如果指定的 Kafka Topic 已存在,可以不用开启 `Describe` 和 `Create` 权限。 | ||||||
|
|
||||||
|
|
@@ -387,9 +381,24 @@ write-key-threshold = 30000 | |||||
| SELECT COUNT(*) FROM INFORMATION_SCHEMA.TIKV_REGION_STATUS WHERE DB_NAME="database1" AND TABLE_NAME="table1" AND IS_INDEX=0; | ||||||
| ``` | ||||||
|
|
||||||
| ## Kafka 消息大小限制 | ||||||
|
|
||||||
| Kafka 会限制每个 Topic 可以接收的消息大小。目标 Topic 当前生效的限制由以下配置决定: | ||||||
|
|
||||||
| | 参数 | 作用 | | ||||||
| | --- | --- | | ||||||
| | Kafka Topic [`max.message.bytes`](https://kafka.apache.org/43/configuration/topic-configs/#topicconfigs_max.message.bytes) | 为指定 Topic 设置消息大小限制,覆盖 broker 的默认值。 | | ||||||
| | Kafka broker [`message.max.bytes`](https://kafka.apache.org/43/configuration/broker-configs/#brokerconfigs_message.max.bytes) | Topic 未设置 `max.message.bytes` 时使用的默认值。 | | ||||||
|
|
||||||
| Kafka Sink 启动时会读取目标 Topic 当前生效的消息大小限制,用于在发送前判断消息是否过大。如果编码后的单条消息超过该限制且未配置大消息处理,Kafka Sink 会返回 `ErrMessageTooLarge`。处理方法参见[Kafka Sink 返回 `ErrMessageTooLarge` 时,如何处理?](/ticdc/troubleshoot-ticdc.md#kafka-sink-返回-errmessagetoolarge-时如何处理)。 | ||||||
|
|
||||||
| > **注意:** | ||||||
| > | ||||||
| > 如果 Kafka Sink 因 Kafka ACL 等原因无法读取 Topic 或 broker 的消息大小配置,会使用 changefeed 的 `max-message-bytes` 作为本地消息大小限制。如果该值与 Kafka 实际生效的限制不一致,Kafka Sink 可能无法准确判断消息能否发送,调大 Kafka 的限制后也可能无法自动恢复。请确保 changefeed 使用的 Kafka 账号具有读取 Topic 和 broker 配置的权限;如果无法授予该权限,需要设置 changefeed 的 `max-message-bytes` 与 Kafka 的消息大小限制一致。 | ||||||
|
|
||||||
| ## 处理超过 Kafka Topic 限制的消息 | ||||||
|
|
||||||
| Kafka Topic 对可以接收的消息大小有限制,该限制由 [`max.message.bytes`](https://kafka.apache.org/documentation/#topicconfigs_max.message.bytes) 参数控制。当 TiCDC Kafka sink 在发送数据时,如果发现数据大小超过了该限制,会导致 changefeed 报错,无法继续同步数据。为了解决这个问题,TiCDC 新增一个参数 `large-message-handle-option` 并提供如下解决方案。 | ||||||
| 当消息超过 Kafka 大小限制时,可以配置 `large-message-handle-option`,避免消息因过大而无法发送。 | ||||||
|
|
||||||
| 目前,如下功能支持 Canal-JSON 和 Open Protocol 两种编码协议。使用 Canal-JSON 协议时,你需要在 `sink-uri` 中设置 `enable-tidb-extension=true`。 | ||||||
|
|
||||||
|
|
||||||
Uh oh!
There was an error while loading. Please reload this page.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.