拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Redpanda Connect PostgreSQL CDC 类型系统全解:从 WAL 解码到 Common Schema 的类型映射指南

Redpanda Connect PostgreSQL CDC 类型系统全解从 WAL 解码到 Common Schema 的类型映射指南【免费下载链接】connectFancy stream processing made operationally mundane项目地址: https://gitcode.com/GitHub_Trending/con/connectpostgres_cdc是 Redpanda Connect 中基于 PostgreSQL 逻辑复制Logical Replication实现的变更数据捕获CDC输入组件。本文以仓库内internal/impl/postgresql/TYPES.md为核心结合pglogicalstream子包源码系统讲解行数据如何从 PostgreSQL 原生类型被解码为 Go 类型、再映射到 Benthos Common Schema 的完整链路帮助你在使用parquet_encode等下游处理器时准确预判数据类型并为排查类型不匹配问题提供源码级依据。整体架构两条独立的数据解码路径postgres_cdc输入通过SetStructuredMut以原生 Go 类型交付行数据。下游消费者有两种取数方式得到的结果在类型语义上保持一致调用AsStructured()例如parquet_encode处理器会直接拿到强类型值调用AsBytes()会得到惰性lazy序列化的 JSON 字符串。产生行数据的两条代码路径完全独立路径驱动方式核心函数特点CDC增量变更pgx v5 解码 WAL 逻辑复制消息decodeTextColumnData针对 pgtype 名称做归一化 switch把 Go 类型修正为与声明 schema 类型一致如 int16 → int32、pgtype.Numeric→ stringSnapshot存量快照标准database/sql扫描prepareScannersAndGetters每个列类型映射到特定的sql.Null*scanner直接产出匹配的 Go 类型两条路径对同一 PostgreSQL 列必须产生完全相同的 Go 类型。该类型集合作为消息元数据中的 schemaBenthos Common Schema 格式暴露给下游因此下游处理器可以放心依赖这些类型无需再做运行时猜测。入口配置与消息装配见 internal/impl/postgresql/input_pg_stream.go。类型映射总表下表是 TYPES.md 给出的权威映射同时覆盖 CDC 与 Snapshot 两条路径PG TypeSchema TypeCDC Go TypeSnapshot Go TypeBOOLBooleanboolboolSMALLINT (int2)Int32int32int32INTEGER (int4)Int32int32int32BIGINT (int8)Int64int64int64REAL (float4)Float32float32float32DOUBLE PRECISION (float8)Float64float64float64NUMERIC / DECIMALStringstringstringTEXT / VARCHAR / CHARStringstringstringBYTEAByteArray[]byte[]byteDATETimestamptime.Timetime.TimeTIMEStringstringstringTIMETZStringstringstringTIMESTAMPTimestamptime.Timetime.TimeTIMESTAMPTZTimestamptime.Timetime.TimeUUIDStringstringstringJSON / JSONBAny(native)(native)注意此处 NUMERIC/DECIMAL 列在无精度/刻度信息时按String兜底当atttypmod或sql.ColumnType.DecimalSize()携带了 precision/scale 信息时schema 会升级为Decimal/BigDecimal详见下文NUMERIC 的精度感知映射。类型映射的源码级实现pgTypeNameToCommonType位于 pglogicalstream/schema.go是映射的核心它把 PostgreSQL 类型名大小写不敏感database/sql返回大写也会被统一处理转换为bschema.CommonType。关键分支包括int2/smallint、int4/integer都收敛为Int32numeric/decimal、text/varchar/character varying/bpchar/name收敛为Stringtime/timetz含time without time zone等变体一律为String原始文本json/jsonb为Any未知类型兜底为Any。该函数与 schema_test.go 中的TestPgTypeNameToCommonType表驱动用例一一对应测试还验证了大写输入如BOOL、INT4与未知类型如INET、_INT4兜底到Any的行为。NUMERIC 的精度感知映射虽然pgTypeNameToCommonType对 NUMERIC 默认返回String但当列声明了精度/刻度时schema 会升级为真正的Decimal/BigDecimal。其实现依赖 PostgreSQL 的atttypmod编码规则PostgreSQL 使用-1表示无修饰符有修饰符时编码为((precision 16) | scale) VARHDRSZ其中VARHDRSZ 4pgNumericModFromAtttypmod解码出 precision/scale小于 4 的 atttypmod 视为缺失同时防御 Go 零值pgNumericToCommon在有修饰符时构造schema.Decimal(precision, scale)否则构造BigDecimal。TestPgNumericModFromAtttypmod覆盖了decimal(18,4)、边界值decimal(1,0)、最大值decimal(38,38)以及无修饰符哨兵值-1等场景可直接作为理解该编码的参考用例。TIMETZ 的 OID 兜底timetzOID 1266是一个值得注意的特例它不在 pgx 默认的类型映射表中pgtype.NewMap()未注册。为此CDC 路径decodeTextColumnData在TypeForOID查不到时走string(data)兜底直接返回原始文本Snapshot 路径database/sql返回的是数字 OID 字符串通过resolveTypeName配合pgOIDToTypeName映射表把1266解析回大写类型名TIMETZ再进入常规映射流程。TestRelationMessageToSchemaTimetz与TestResolveTypeName分别验证了这两条支路的正确性。需要注意的是resolveTypeName只对已知 OID 做解析未知的数字 OID 字符串会原样透传。CDC 路径WAL 消息解码详解增量变更走的是 pgx v5 的逻辑复制解码链路核心函数decodeTextColumnData位于 pglogicalstream/replication_message_decoders.go。它对每个列值执行以下归一化逻辑分支行为uuidpgx 解码为[16]uint8转换为uuid.UUID(...).String()字符串tsrange返回经sanitizeTsrange去除引号后的文本兼容旧 pgtype 行为见 pgtype_compat.goint2pgx 解码为int16提升为int32以匹配 Int32 schema 类型numeric返回规范化十进制字符串详见下文NaN/Infinity/-Infinity原样透传datetime.Time±infinity 无法表示时返回niltime返回原始文本字符串刻意绕开pgtype.Time结构体tsvectorpgx 解码为pgtype.TSVector结构体此处返回原始文本字符串timestamp/timestamptztime.Time±infinity 返回nil其他返回 pgx 的默认解码值在更上层toStreamMessage负责把 pgoutput 格式的逻辑复制消息Begin/Insert/Update/Delete等转换为统一的StreamMessageUpdateMessage在REPLICA IDENTITY FULL场景下会用旧元组OldTuple作为 TOAST 未变化列的兜底值来源decodeTuple则逐列处理null、unchanged toast、text三种 tuple 数据类型。这些解码结果随后被装配为带元数据的消息见input_pg_stream.go的processStreamtable、operation、lsn、schema、commit_ts_ms、before。NUMERIC 的规范化输出numeric分支的规范化实现位于 internal/sqlutil/decimal.go当atttypmod携带 precision/scale 时调用CanonicaliseDecimal(text, precision, scale)输出精确 scale 位小数不足右补零、无前导零保留小数点前单个 0、不使用科学计数法、负数仅带前导-否则调用CanonicaliseBigDecimal(text)规范化自然刻度形式超出列精度或无法在声明刻度下精确表示的值会被拒绝并报错而不是静默四舍五入——这一点对金融类等精度敏感场景至关重要。NaN、Infinity、-Infinity由于没有规范的十进制表示直接原样透传。此外WAL 行数据如果无法被 JSON 序列化实践中主要是非有限浮点值 NaN/Infinity消息会被打上 error 标记并以纯文本形式发布配合errored()Bloblang 函数和switch输出等错误处理组件即可路由到死信队列而复制检查点仍会正常推进详见input_pg_stream.go的processStream实现。Snapshot 路径database/sql 扫描详解存量快照由 pglogicalstream/snapshotter.go 中的prepareScannersAndGetters完成。它基于快照查询返回的sql.ColumnType为每列分配 scanner 与 getter数据库类型名Scanner产出 Go 类型VARCHAR / TEXT / UUIDsql.NullStringstringBOOLsql.NullBoolboolINT2 / INT4sql.NullInt32int32INT8sql.NullInt64int64FLOAT4sql.NullFloat64收窄float32FLOAT8sql.NullFloat64float64DATE / TIMESTAMP / TIMESTAMPTZsql.NullTimetime.TimeJSON / JSONBsql.NullStringjson.Unmarshalmap[string]any/[]any/float64/string/bool/nil树INETsql.NullStringnetip解析CIDR 前缀字符串裸地址补全主机位长度TSRANGEsql.NullStringsanitizeTsrange去引号文本_INT4 / _TEXT数组sql.NullStringpgTypeMap.SQLScanner[]*int32/[]*stringNULL 元素保留为 JSON null其他含 NUMERICsql.NullStringstring值得注意的实现细节JSON/JSONB 经过json.Unmarshal展开为 stdlib 类型树任何sql.*包装器都不会泄漏到下游FLOAT4 通过sql.NullFloat64扫描后收窄为float32与 CDC 路径从 pgx 原生拿到的float32保持一致——这正是两条路径必须产生相同 Go 类型约束的具体落地。Schema 元数据下游处理器如何获得类型列类型信息以 Benthos Common Schema 形式随消息元数据下发由两个函数构造CDC 路径relationMessageToSchemaschema.go——基于RelationMessage的列 OID 与TypeModifier通过 pgx 类型映射未知 OID 用pgOIDToTypeName兜底逐列构造Common所有列默认Optional: true外层为以表名为名的Object类型Snapshot 路径columnTypesToSchema——基于sql.ColumnType切片对 NUMERIC 读取DecimalSize()得到 precision/scale其余同 CDC 一致。TestRelationMessageToSchema 演示了一个含 9 列的orders表amount列在无修饰符TypeModifier: -1时被断言为BigDecimal而TestRelationMessageToSchemaDecimalWithModifier中numeric(18,4)列被断言为Decimal且Precision18、Scale4。运行时输入组件通过MetaSetImmut(schema, ...)把该 schema 挂到每条消息上input_pg_stream.goparquet_encode等处理器即可据此直接消费强类型值。使用建议与注意事项精度敏感数据优先用 NUMERIC 并声明精度未声明精度的 NUMERIC 在 schema 中体现为BigDecimal声明的则体现为Decimal(precision, scale)后者输出按 scale 补齐、可预测性更强时间边界值会变nilDATE/TIMESTAMP/TIMESTAMPTZ 的 ±infinity 无法表示为time.Time两条路径都返回nil下游需做好空值处理TIME/TIMETZ 是原始文本而非时间对象如需时间运算应使用 Bloblang 的时间解析方法做显式转换未知类型退化为Any/字符串自定义类型、INET、数组等不在映射表中的类型schema 类型为Any但 Snapshot 路径对 INET、_INT4、_TEXT等做了专门的文本/数组解码实际 Go 类型以文档表格与源码为准TOAST 未变化列当REPLICA IDENTITY非 FULL 时UPDATE/DELETE 中未变化的 TOAST 列会以unchanged_toast_value配置项指定的占位值默认null呈现可在输入配置中显式设置哨兵字符串以便下游识别。关键文件索引internal/impl/postgresql/TYPES.md — 类型系统权威文档本文核心来源internal/impl/postgresql/pglogicalstream/schema.go — 类型名映射、OID 兜底、schema 构造internal/impl/postgresql/pglogicalstream/replication_message_decoders.go — CDC 解码与类型归一化internal/impl/postgresql/pglogicalstream/snapshotter.go — 快照扫描与 scanner 映射internal/impl/postgresql/pglogicalstream/schema_test.go — 类型映射与 schema 构造的测试佐证internal/sqlutil/decimal.go — NUMERIC 规范化输出实现internal/impl/postgresql/input_pg_stream.go —postgres_cdc输入配置与消息元数据装配【免费下载链接】connectFancy stream processing made operationally mundane项目地址: https://gitcode.com/GitHub_Trending/con/connect创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门