diff --git a/ticdc/ticdc-faq.md b/ticdc/ticdc-faq.md index 359473732894..6d2d086a3c8c 100644 --- a/ticdc/ticdc-faq.md +++ b/ticdc/ticdc-faq.md @@ -238,12 +238,6 @@ cdc cli changefeed create --server=http://127.0.0.1:8300 --sink-uri="kafka://127 * `replica.fetch.max.bytes`,将 Kafka 的 `server.properties` 中该参数调大到 `1073741824` (1 GB)。 * `fetch.message.max.bytes`,适当调大 `consumer.properties` 中该参数,确保大于 `message.max.bytes`。 -## TiCDC 把数据同步到 Kafka 时,能在 TiDB 中控制单条消息大小的上限吗? - -对于 Avro 和 Canal-JSON 格式,消息是以行变更为单位发送的,一条 Kafka Message 仅包含一条行变更。一般情况下,消息的大小不会超过 Kafka 单条消息上限,因此,一般不需要限制单条消息大小。如果单条 Kafka 消息大小确实超过 Kafka 上限,请参考[为什么 TiCDC 到 Kafka 的同步任务延时越来越大](/ticdc/ticdc-faq.md#为什么-ticdc-到-kafka-的同步任务延时越来越大)。 - -对于 Open Protocol 格式,一条 Kafka Message 可能包含多条行变更。因此,有可能存在某条 Kafka Message 消息过大。可以通过 `max-message-bytes` 控制每次向 Kafka broker 发送消息的最大数据量(可选,默认值 10 MB),通过 `max-batch-size` 参数指定每条 kafka 消息中变更记录的最大数量(可选,默认值 `16`)。 - ## 在一个事务中对一行进行多次修改,TiCDC 会输出多条行变更事件吗? 不会,在进行事务操作时,对于在一个事务内多次修改同一行的情况,TiDB 仅会将最新一次的修改结果发送给 TiKV。因此 TiCDC 仅能获取到最新一次修改的结果。 diff --git a/ticdc/ticdc-open-protocol.md b/ticdc/ticdc-open-protocol.md index d1f6d8e065ea..67fbea7d8f8c 100644 --- a/ticdc/ticdc-open-protocol.md +++ b/ticdc/ticdc-open-protocol.md @@ -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-消息大小限制)。 + +当 Message 中的 Row Changed Event 数量已经达到 `max-batch-size`,或者继续加入下一个 Row Changed Event 会使 Message 超过实际生效的大小阈值时,后续 Row Changed Event 会被写入一条新的 Message。 + ## Event 格式定义 本部分介绍 Row Changed Event、DDL Event 和 Resolved Event 的格式定义。 diff --git a/ticdc/ticdc-sink-to-kafka.md b/ticdc/ticdc-sink-to-kafka.md index 27591c2f0ca5..3a4c477e5371 100644 --- a/ticdc/ticdc-sink-to-kafka.md +++ b/ticdc/ticdc-sink-to-kafka.md @@ -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-数量和大小)。| +| `max-batch-size` | Open Protocol Message 最多包含的 Row Changed Event 数量(可选,默认值 `16`)。详细说明参见[控制 Message 中的 Event 数量和大小](/ticdc/ticdc-open-protocol.md#控制-message-中的-event-数量和大小)。| | `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`。 diff --git a/ticdc/troubleshoot-ticdc.md b/ticdc/troubleshoot-ticdc.md index cd39eb9df321..be59c7b2a333 100644 --- a/ticdc/troubleshoot-ticdc.md +++ b/ticdc/troubleshoot-ticdc.md @@ -102,19 +102,35 @@ Warning: Unable to load '/usr/share/zoneinfo/zone1970.tab' as time zone. Skippin - 如果 PD 是由 v4.0.8 或更低版本滚动升级到新版,详见 [PD issue #3366](https://github.com/tikv/pd/issues/3366)。 - 对于其他情况,请将上述命令执行结果反馈到 [AskTUG 论坛](https://pingkai.cn/tidbcommunity/forum/tags/ticdc)。 -## 使用 TiCDC 同步消息到 Kafka 时 Kafka 报错 `Message was too large`,该如何处理? +## Kafka Sink 返回 `ErrMessageTooLarge` 时,如何处理? -仅在 Sink URI 中为 Kafka 配置 `max-message-bytes` 参数不能有效控制输出到 Kafka 的消息大小,需要在 Kafka server 配置中加入如下配置以增加 Kafka 接收消息的字节数限制。 +处理方式取决于 TiCDC 版本: + +**v8.5.8 及之后版本**:根据错误信息中的消息大小,调大 Kafka Topic 的 `max.message.bytes`。Kafka Sink 重建时会重新读取该配置,changefeed 随后自动恢复同步。 + +如果 changefeed 未自动恢复,请确认其使用的 Kafka 账号具有读取 Topic 和 broker 配置的权限。更多信息参见[Kafka 消息大小限制](/ticdc/ticdc-sink-to-kafka.md#kafka-消息大小限制)。 + +**v8.5.8 之前的版本**: + +1. 根据错误信息确认待发送消息的大小。 +2. 将 Kafka Topic 的 `max.message.bytes` 调整为不小于该消息大小。 +3. 暂停 changefeed,将 `max-message-bytes` 设置为与 `max.message.bytes` 相同的值,然后恢复 changefeed。 + +调整 Kafka 消息大小限制时,还需要检查以下相关配置: ``` -# broker 能接收消息的最大字节数 -message.max.bytes=2147483648 +# Topic 能接收消息的最大字节数 +max.message.bytes=<不小于待发送消息的大小> +# broker 能接收消息的最大字节数,使用 broker 默认限制时配置 +message.max.bytes=<不小于待发送消息的大小> # broker 可复制的消息的最大字节数 -replica.fetch.max.bytes=2147483648 +replica.fetch.max.bytes=<不小于 message.max.bytes> # 消费者端的可读取的最大消息字节数 -fetch.message.max.bytes=2147483648 +fetch.message.max.bytes=<不小于 Kafka 中实际允许的消息大小> ``` +如果不希望调大 Kafka 的消息大小限制,可以配置 `large-message-handle-option`,使用 Claim-Check 或只输出 Handle Key 的方式处理大消息。更多信息参见[处理超过 Kafka Topic 限制的消息](/ticdc/ticdc-sink-to-kafka.md#处理超过-kafka-topic-限制的消息)。 + ## TiCDC 同步时,在下游执行 DDL 语句失败会有什么表现,如何恢复? 如果某条 DDL 语句执行失败,同步任务 (changefeed) 会自动停止,checkpoint-ts 断点时间戳为该条出错 DDL 语句的结束时间戳 (finish-ts)。如果希望让 TiCDC 在下游重试执行这条 DDL 语句,可以使用 `cdc cli changefeed resume` 恢复同步任务。例如: @@ -130,7 +146,7 @@ cdc cli changefeed resume -c test-cf --server=http://127.0.0.1:8300 3. 恢复被暂停的 changefeed。 > **注意:** -> +> > 虽然将 changefeed 的 `start-ts` 设为报错时的 `checkpoint-ts` 值加上 1,然后重建任务也可以跳过该 DDL 语句,但同时会导致 TiCDC 丢失 `checkpointTs+1` 时刻对应的 DML 数据变更。严禁在生产环境执行这样的操作。 ```shell @@ -143,5 +159,5 @@ cdc cli changefeed create --server=http://127.0.0.1:8300 --sink-uri="mysql://roo 该问题通常是由于 TiCDC 与 Kafka 集群连接失败导致。你可以通过检查 Kafka 的日志以及网络状况来排查。一个常见的原因是在创建同步任务时没有指定正确的 `kafka-version` 参数,导致 TiCDC 内部的 Kafka client 在访问 Kafka server 时使用了错误的 Kafka API 版本。你可以通过配置 [`--sink-uri`](/ticdc/ticdc-sink-to-kafka.md#sink-uri-配置-kafka) 指定正确的 `kafka-version` 参数来修复。例如: ```shell -cdc cli changefeed create --server=http://127.0.0.1:8300 --sink-uri "kafka://127.0.0.1:9092/test?topic=test&protocol=open-protocol&kafka-version=2.4.0" +cdc cli changefeed create --server=http://127.0.0.1:8300 --sink-uri "kafka://127.0.0.1:9092/test?topic=test&protocol=open-protocol&kafka-version=2.4.0" ```