ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

MySQL CDC 类型系统深度解析:Redpanda Connect 的 MySQL 列类型到 Go/Schema 类型映射全指南

MySQL CDC 类型系统深度解析:Redpanda Connect 的 MySQL 列类型到 Go/Schema 类型映射全指南 MySQL CDC 类型系统深度解析Redpanda Connect 的 MySQL 列类型到 Go/Schema 类型映射全指南【免费下载链接】connectFancy stream processing made operationally mundane项目地址: https://gitcode.com/GitHub_Trending/con/connect本文围绕 Redpanda Connect 的mysql_cdc输入组件位于 internal/impl/mysql系统讲解其内置类型系统Type System的完整设计CDC 增量流与 Snapshot 快照两条数据生产路径如何产生一致的 Go 类型、MySQL 各列类型在 Schema 元数据中的映射规则以及这些类型如何被parquet_encode等下游处理器消费。读完本文你将掌握 mysql_cdc 的底层类型映射原理能够准确预测任意 MySQL 表结构在消息中呈现的 Go 类型与 Schema 声明并据此正确配置下游处理器。一、类型系统的设计目标与总体架构mysql_cdc输入通过SetStructuredMut将行数据以原生 Go 类型投递给下游调用AsStructured()的消费者例如parquet_encode处理器会直接收到强类型值调用AsBytes()的消费者则会获得惰性lazy序列化的 JSON。这一设计意味着类型正确性必须在输入端就得到保证而不是依赖下游做二次推断。为此仓库内存在两条相互独立、但必须产出完全相同 Go 类型的数据生产路径对应文档 TYPES.md 的 Overview路径底层库归一化入口CDC增量流go-mysql 的 canal 库解码 binlog 事件mapMessageColumninput_mysql_stream.goSnapshot快照标准库database/sql扫描prepSnapshotScannerAndMappersinput_mysql_stream.go两条路径的一致性约束是系统的核心不变量同一 MySQL 列无论来自增量事件还是快照扫描都必须产出相同的 Go 类型。而表 Schema以消息元数据schema的形式暴露则忠实地反映这些类型供下游处理器如parquet_encode依赖。二、完整类型映射表以下映射表直接来自 TYPES.md列出了 MySQL 列类型到 Schema 类型、CDC Go 类型、Snapshot Go 类型的完整对应关系MySQL TypeSchema TypeCDC Go TypeSnapshot Go TypeTINYINTInt32int32int32SMALLINTInt32int32int32MEDIUMINTInt32int32int32INTInt32int32int32UNSIGNED TINYINTInt32int32int32UNSIGNED SMALLINTInt32int32int32UNSIGNED MEDIUMINTInt32int32int32UNSIGNED INTInt64int64int64BIGINTInt64int64int64UNSIGNED BIGINTInt64int64int64YEARInt32int32int32FLOATFloat32float32float32DOUBLEFloat64float64float64DECIMAL / NUMERICStringstringstringDATETimestamptime.Timetime.TimeDATETIMETimestamptime.Timetime.TimeTIMESTAMPTimestamptime.Timetime.TimeTIMEStringstringstringBITInt64int64int64CHAR / VARCHAR / TEXTStringstringstringBINARY / VARBINARY / BLOBByteArray[]byte[]byteENUMStringstringstringSETArray[String][]any[]anyJSONAny(native)(native)注意上表中 DECIMAL / NUMERIC 在 Schema 层的“String”仅为文档中的简化表述从源码看mysqlColumnToCommon对 DECIMAL 实际会构造schema.Decimal逻辑类型携带 precision/scale见下文第三节这是 Schema 元数据层面比普通字符串更精确的表示。2.1 关键设计注记Notes文档 TYPES.md 中列出的设计注记对应到源码均有明确实现整数宽度BIGINT 与 UNSIGNED INT 之所以使用 Int64是因为它们的最大值超出 int32 范围其余整数类型均能装进 int32。在 schema.go 中可见mysqlColumnToCommon对TYPE_NUMBER的判断逻辑当原始类型以bigint开头或int开头且IsUnsigned为真时映射为schema.Int64。DECIMAL 用字符串表示以保留任意精度若改用 float64 会静默丢失数字精度。JSON两条路径都会执行json.Unmarshal产出一棵标准库类型树map[string]any、[]any、float64、string、bool、nil不会有任何sql.*包装类型泄漏出去。零值日期时间Zero datetimesCDC 路径把非法日期时间例如0000-00-00 00:00:00以字符串形式投递mapMessageColumn将其转换为nil见下文 2.2。UNSIGNED BIGINT 超出 MaxInt64当值超过math.MaxInt64时以uint64原样透传这是一个绝大多数下游消费者不会遇到的边界情形。2.2 零值日期时间的处理细节在 input_mysql_stream.go 的TYPE_DATETIME/TYPE_TIMESTAMP分支中若值类型是string即 canal 以字符串形式给出的非法日期mapMessageColumn直接返回nil若已是time.Time则原样返回。对应测试TestMapMessageColumn中的 zero datetime string to nil 用例schema_test.go验证了这一行为。而 Snapshot 路径中sql.NullTime在非法日期下Valid为 false同样由 mapper 映射为nilinput_mysql_stream.go。三、Schema 生成从 MySQL 表结构到 Benthos Common Schemamysql_cdc除了投递行数据还会把每个被追踪表的 Schema 以 Benthos Common Schema 格式写入消息元数据schema该行为在 mysql_cdc.adoc 中有说明兼容parquet_encode等处理器。这一转换的核心函数是mysqlColumnToCommon与mysqlTableToCommonSchema位于 schema.go。3.1 转换流程// 伪代码schema.go 中的核心逻辑 mysqlTableToCommonSchema(table): children [mysqlColumnToCommon(col) for col in table.Columns] // 保持列顺序 return Common{ Name: table.Name, Type: Object, Children: children } mysqlColumnToCommon(col): switch col.Type: TYPE_NUMBER - bigint 或 unsigned int Int64否则 Int32 TYPE_MEDIUM_INT - Int32 TYPE_FLOAT - double 开头 Float64否则 Float32 TYPE_DECIMAL - schema.NewDecimal(name, precision, scale)解析自 RawType TYPE_STRING - String TYPE_DATETIME/TIMESTAMP/TYPE_DATE - Timestamp TYPE_TIME - String TYPE_BINARY - ByteArray TYPE_BIT - Int64 TYPE_ENUM - String TYPE_SET - Array[String]Children 含 element: String TYPE_JSON - Any TYPE_POINT - ByteArray几何类型暂按二进制处理 default - 报错 unsupported MySQL column type几点值得注意的细节列顺序被严格保留mysqlTableToCommonSchema按table.Columns的顺序构建 Children测试TestMysqlTableToCommonSchemaschema_test.go显式断言了列顺序被保留。所有列默认 OptionalmysqlColumnToCommon返回的Common.Optional恒为true注释说明All MySQL columns can be NULL unless specified otherwise测试中也统一断言result.Optional trueschema_test.go。DECIMAL 携带 Logical 信息parseMySQLDecimal用正则解析RawType如decimal(18,4)、decimal(7)、裸decimal、带unsigned后缀等形式给出 MySQL 默认值裸decimal视为(10, 0)、decimal(p)视为(p, 0)。解析成功后schema.NewDecimal会把 precision/scale 存入Logical.Decimal测试TestMysqlColumnToCommonDecimalCarriesLogicalschema_test.go验证了这一点解析失败如varchar(20)、空串则直接报错。JSON 列声明为Any源码注释schema.go明确说明Any向parquet_encode等下游消费者传达字段类型未知的信号消费者必须显式处理Any否则会返回可操作的报错提示用户添加类型转换步骤。3.2 Schema 的缓存与失效schema.go之外Schema 的获取、缓存与失效位于 input_mysql_stream.gogetTableSchema按表名缓存已序列化的 SchemacommonSchema.ToAny()getOrExtractTableSchemaByName快照消息可能还没有 canal Table 对象此时返回 nil等后续 CDC 事件到达时再补上invalidateTableSchema在OnTableChangedDDL 变更建表、改表、重命名、删表时使缓存失效保证 Schema 与最新表结构一致序列化后的 Schema 支持schema.ParseFromAny无损还原测试TestMysqlTableToCommonSchemaRoundtripschema_test.go验证了往返一致性。四、CDC 路径的类型归一化mapMessageColumnCDC 路径由 go-mysql 的 canal 库解码 binlog 事件。canal 给出的 Go 值类型并不总是与声明的 Schema 类型一致因此需要mapMessageColumninput_mysql_stream.go做归一化。文档中 int8 → int32 的例子正是指这种提升。4.1 各分支归一化规则列类型输入类型输出TYPE_NUMBERint / uint / int64 / uint64 等int8/int16/uint8/uint16 →int32int32 透传int / int64 / uint / uint32 →int64uint64 若 math.MaxInt64保持uint64否则转 int64TYPE_MEDIUM_INTint32 / uint32int32uint32 直接转 int32TYPE_FLOAT任意原样透传canal 已给出 float32/float64TYPE_DECIMALstring解析 precision/scale 后经sqlutil.CanonicaliseDecimal规范化返回规范化的字符串TYPE_SETint64位图按col.SetValues逐位解码为[]any每个成员一个元素TYPE_DATEstring / time.Timestring 用time.Parse(2006-01-02, d)转为 time.Timetime.Time 透传TYPE_DATETIME / TYPE_TIMESTAMPstring返回nil零值日期处理TYPE_ENUMint64序号校验 1 ≤ ordinal ≤ len(EnumValues) 后返回对应字符串越界报错TYPE_JSONstringjson.Unmarshal解码为标准库类型树TYPE_STRINGstring / []byte非 blob 列返回 stringblob 类RawType 含 blobfallthrough 到 BINARY 分支TYPE_BINARY[]byte / string返回 []byte4.2 测试验证TestMapMessageColumnschema_test.go以表格驱动方式覆盖了上述核心场景例如int8(42) → int32(42)、int16(1000) → int32(1000)、int32/int64透传uint8(255) → int32(255)、uint16(65535) → int32(65535)、uint32(4294967295) → int64(4294967295)uint64(math.MaxInt64 1)保持uint64超长 DECIMAL 字符串999999999999999999999999999999999999.99原样透传DATE 字符串2024-12-10→time.Time零值 DATETIME 字符串 →nilnil值直接透传MySQL NULL 语义。五、Snapshot 路径的类型扫描prepSnapshotScannerAndMappers快照路径不使用 canal而是标准database/sql扫描。prepSnapshotScannerAndMappersinput_mysql_stream.go为每一列选择特定的sql.Null*扫描器使产出的 Go 类型与 CDC 路径完全一致。5.1 扫描器选择表DatabaseTypeName扫描器产出类型BINARY / VARBINARY / TINYBLOB / BLOB / MEDIUMBLOB / LONGBLOBsql.Null[[]byte][]byteDATETIME / TIMESTAMPsql.NullTimetime.TimeTINYINT / SMALLINT / MEDIUMINT / INT / YEAR / UNSIGNED TINYINT / UNSIGNED SMALLINT / UNSIGNED MEDIUMINTsql.NullInt32int32BIGINT / UNSIGNED INT / UNSIGNED BIGINTsql.NullInt64int64DECIMAL / NUMERICsql.NullStringCanonicaliseDecimalstringFLOATsql.Null[float32]float32DOUBLEsql.Null[float64]float64SETsql.NullString按逗号拆分为[]any[]anyJSONsql.NullStringjson.Unmarshal标准库类型树BITsql.Null[[]byte]按字节拼接为 int64int64DATEsql.NullTimetime.Time其他/默认sql.Null[string]string所有 mapper 在Valid falseSQL NULL时返回nil与 CDC 路径对 NULL 的nil透传语义对齐。5.2 快照扫描的工程细节从 input_mysql_stream.go 的snapshotTable可以看到快照按主键排序、以snapshot_max_batch_size默认 1000分页拉取WHERE (pk...) (lastSeen...)的 keyset 分页保证一致性表必须有主键否则报错snapshot.go快照行以MessageOperationReadoperation 元数据为read见 event.go投递快照阶段还会统计mysql_snapshot_rows_processed_total指标按表维度计数。六、两种典型消费场景6.1 场景一类型化的下游处理器如 parquet_encodemysql_cdc的 Schema 元数据被设计为与parquet_encode兼容。消息携带的schema元数据Benthos Common Schema 格式可让parquet_encode直接获得字段类型。需要注意DECIMAL 列在 Schema 中携带 precision/scale 逻辑类型这为 parquet 的 decimal 物理类型提供了必要信息JSON 列声明为Anyparquet_encode无法自动推断其结构需要显式添加类型转换处理器否则会返回可操作的报错提示。6.2 场景二消息元数据驱动的路由与过滤mysql_cdc为每条消息写入以下元数据字段见 mysql_cdc.adoc 及 input_mysql_stream.go 的MetaSet调用operationinsert / update / delete / readread 对应快照消息table表名binlog_positionbinlog 位置仅 CDC 消息快照消息不设置schemaBenthos Common Schema 格式的表结构。七、完整配置示例以下完整配置含各字段默认值摘自 mysql_cdc.adocinput: label: mysql_cdc: flavor: mysql # mysql | mariadb dsn: user:passwordtcp(localhost:3306)/database # 必填 tables: [] # 必填至少一个表 checkpoint_cache: # 必填存储 binlog 断点的 cache 资源 checkpoint_key: mysql_binlog_position snapshot_max_batch_size: 1000 stream_snapshot: false # 是否先全量快照 max_parallel_snapshot_tables: 1 auto_replay_nacks: true checkpoint_limit: 1024 max_reconnect_attempts: 10 # 高级字段 batching: count: 0 byte_size: 0 period: check: 补充说明几个与类型系统直接相关的字段配置定义见 input_mysql_stream.goflavor连接 MySQL 或 MariaDB 风格数据库gomysql.MySQLFlavor/gomysql.MariaDBFlavorstream_snapshot为true时连接器先做全量快照再进入增量流为false时从当前 binlog 位置开始checkpoint_cache/checkpoint_keybinlog 位置断点存储重启后从断点续传而非全量重读多路 mysql_cdc 共享 cache 时可换 keymax_parallel_snapshot_tables并行快照的表数量默认 1snapshot_max_batch_size快照单批最大行数默认 1000checkpoint_limit同一时刻可处理的最大消息数提升该值可启用并行处理与输出端批处理但为保证至少一次投递同一 binlog 偏移下的消息必须按序确认。八、快照一致性机制补充快照路径snapshot.go采用FLUSH TABLES WITH READ LOCKSTART TRANSACTION WITH CONSISTENT SNAPSHOT的组合来保证一致性用锁连接对目标表执行FLUSH TABLES ... WITH READ LOCK阻止写入在锁窗口内为每个 worker 开启START TRANSACTION WITH CONSISTENT SNAPSHOT只读事务repeatable read 隔离级别使所有 worker 读到完全一致的数据视图一个 worker 事务可复用于多张表仍在锁窗口内读取 binlog 位置MySQL 8.4 用SHOW BINARY LOG STATUS旧版本回退到SHOW MASTER STATUS见 snapshot.go立即释放锁避免长时间阻塞其他连接。快照完成后会发送内部哨兵事件messageOperationSnapshotComplete将 checkpoint 推进到快照起始的 binlog 位置从而保证快照 增量无缝衔接、不丢不重相关常量的注释见 event.go。九、总结与最佳实践两条路径、一种类型CDC 的mapMessageColumn与快照的prepSnapshotScannerAndMappers必须对同一 MySQL 列产出相同 Go 类型这是本类型系统的核心不变量由 schema_test.go 中的测试用例固化保障。精度安全优先DECIMAL 以规范化字符串表示、JSON 解码为标准库类型树、BINARY/BLOB 保持 []byte这些都是为了防止精度丢失或包装类型泄漏。Schema 元数据是下游的契约消息携带的schema元数据与parquet_encode兼容但 JSON 列的Any类型需要下游显式处理。NULL 语义统一为 nil两条路径对 SQL NULL 与零值日期统一映射为nil下游无需区分来源。边界情况要知晓UNSIGNED BIGINT 超过math.MaxInt64时会以uint64透传快照表必须有主键。若想深入源码建议按以下顺序阅读TYPES.md总览→ schema.go类型映射→ input_mysql_stream.go两路径实现→ schema_test.go测试固化→ snapshot.go快照一致性完整的组件配置文档见 mysql_cdc.adoc。【免费下载链接】connectFancy stream processing made operationally mundane项目地址: https://gitcode.com/GitHub_Trending/con/connect创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表