TiDB 到 Kafka
| 功能 | 说明 |
|---|---|
Schema Migration | If the specified Topic after mapping does not exist in the Target, BladePipe will automatically create the Topic, allowing setting the number of partitions. |
Full Data Migration | Migrate data by sequentially scanning data in tables and writing it in batches to the target database. |
Incremental Data Sync | Sync of common DML like INSERT, UPDATE, DELETE is supported. |
Subscription Modification | Add, delete, or modify the subscribed tables with support for historical data migration. For more information, see Modify Subscription. |
Position Resetting | Reset positions by timestamp to consume again the incremental data that has not been collected as garbage by TiKV in a past period. |
Metadata Retrieval | Retrieve the target metadata with filtering conditions or target primary keys set from the source table. |
高级功能
| 功能 | 说明 |
|---|---|
消息格式 | 支持以下消息格式,文档:消息格式说明
|
Topic Mapping Rules | By default, Topic is formed by connecting source instance id, database, and table with . in between (e.g., my-vgpq6q097174t6t.dingtax.app_key). Also, it supports the mapping rules, namely, keeping the name the same as that in Source, converting the text to lowercase, converting the text to uppercase. |
Table-level Topic | Create Topics corresponding to the tables in the Source, and the table partitions can be obtained automatically. |
DDL Dedicated Topic | Allow specifying a Topic for DDL. If not specified, DDL time is placed in partition 0 of the Topic created from the corresponding table. |
Scheduled Full Data Migration | For more information, see Create Scheduled Full Data DataJob. |
Custom Code | For more information, see Custom Code Processing, Debug Custom Code and Logging in Custom Code. |
Data Filtering Conditions | Support data filtering using WHERE conditions, with SQL-92 as the SQL language. For more information, see Data Filtering. |
使用示例
| 标题 | 详情 |
|---|---|
跨互联网数据互通 (Kafka) | |
Kafka 数据中转校验 | 文档:Kafka 数据中转校验 |
前置条件
| 条件 | 说明 |
|---|---|
账号权限 | 文档:TiDB 需要的权限 |
PD节点网络连通 | 请确保 CloudCanal 各节点能正常与 PD 各节点通讯
|
TiKV GC 回收频率 | 在 TiDB Server 中修改 GC 周期时间为 24小时 以上
|
TiKV 历史变更数据缓存 | 建议根据任务所需适当调整大小
|
任务参数
| 参数名称 | 说明 |
|---|---|
printDetailLog | 打印接收到的增量,常用于判断源端是否有增量数据推送 |
pdHost | 任务请求的 PD 节点地址,格式为: [PD_IP]:[PD_PORT], 多个 PD 节点用 , 隔开 |
cdcGrpcTimeout | 任务与 PD 节点 gRpc 连接通道的超时时间,单位ms |
cdcStubTimeout | gRpc 通道中的每个 stub 的超时时间,超过该时间会自动重新订阅,单位ms |
fastFailKeywords | 字符串数组,以逗号分隔,当异常信息中包含这些关键字时,任务不再尝试重连,直接重启。例如 DEADLINE_EXCEEDED 表示当 gRPC 超时异常时不再重连,直接重启任务 |
Tips: 通用参数配置请参考 通用参数及功能
任务参数
| 参数名称 | 说明 |
|---|---|
schemaFormat | 消息格式,文档:消息格式说明 |
batchWriteSize | 单条消息最大数据条数,超过则拆分消息 |
defaultTopic | 无法找到对应 Topic 的消息则发送到此 Topic (如新增表) |
ddlTopic | 专门发送 DDL 的 Topic, 为空则发送到对应 Topic 的第 0 个分区 |
compressionType | Kafka compression.type 参数, 设置压缩算法, 支持 GZIP, SNAPPY, LZ4, ZSTD 算法 |
batchSize | Kafka batch.size 参数 |
acks | Kafka acks 参数, 默认 all |
maxRequestBytes | Kafka max.request.size 参数 |
lingerMs | Kafka linger.ms 参数, 默认 1 |
envelopSchemaInclude | 当 schemaFormat 设置为 DEBEZIUM_ENVELOP_JSON_FOR_MQ 时,消息体是否包含 schema 信息 |
customClientProps | 自定义传入到 Kafka Client 参数,JSON 格式,key为参数名,value为参数值。此配置项以最高优先级生效。例如:AWS IAM 访问控制 |
Tips: 通用参数配置请参考 通用参数及功能