Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 0 additions & 6 deletions ticdc/ticdc-faq.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 仅能获取到最新一次修改的结果。
Expand Down
13 changes: 13 additions & 0 deletions ticdc/ticdc-open-protocol.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 的格式定义。
Expand Down
39 changes: 24 additions & 15 deletions ticdc/ticdc-sink-to-kafka.md
Original file line number Diff line number Diff line change
Expand Up @@ -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-数量和大小)。|

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

为什么还要保留 max-message-bytes 这个参数?

| `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 实例所需的证书文件路径(可选)。 |
Expand All @@ -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 认证时配置样例如下所示:
Expand Down Expand Up @@ -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` 权限。

Expand Down Expand Up @@ -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`。

Expand Down
32 changes: 24 additions & 8 deletions ticdc/troubleshoot-ticdc.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 中实际允许的消息大小>
Comment thread
coderabbitai[bot] marked this conversation as resolved.
```

如果不希望调大 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` 恢复同步任务。例如:
Expand All @@ -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
Expand All @@ -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"
```
Loading