Apache Pulsar JDBC Sink Connector 完全指南:将 Topic 消息持久化到 ClickHouse / MariaDB / PostgreSQL / SQLite
2026/9/23 10:00:03 网站建设 项目流程
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

导读

本文围绕 site2/docs/io-jdbc-sink.md 展开,系统讲解 Apache Pulsar 的 JDBC sink connector:它负责从 Pulsar topic 拉取消息,并通过 JDBC 持久化到 ClickHouse、MariaDB、PostgreSQL、SQLite 四种数据库。读完本文,你将掌握该连接器的全部配置属性及其默认值、四种数据库的 JSON/YAML 配置写法、基于消息ACTION属性触发 INSERT / UPDATE / DELETE 的写入机制,以及从源码层面理解其连接管理、SQL 自动构建与批量 flush 的实现原理。

当前版本(仓库基线为 2.10.6-SNAPSHOT)中,JDBC sink 支持 INSERT、DELETE 和 UPDATE 三种数据库操作。

一、连接器概览:一条 Topic 与四类数据库之间的桥梁

JDBC sink connector 是 Pulsar IO 连接器家族(位于 pulsar-io/jdbc 目录)中面向关系型数据库的一类实现。与其它 sink 不同,它被拆分为一个公共核心模块与四个数据库专属模块:

数据库仓库模块Sink 类型(sinkType
ClickHousepulsar-io/jdbc/clickhousejdbc-clickhouse
MariaDBpulsar-io/jdbc/mariadbjdbc-mariadb
PostgreSQLpulsar-io/jdbc/postgresjdbc-postgres
SQLitepulsar-io/jdbc/sqlitejdbc-sqlite

从源码看,四个专属模块的实现类本身几乎是空壳——它们只通过@Connector注解声明连接器名称、类型与配置类,实际逻辑全部继承自核心模块。例如:

  • PostgresJdbcAutoSchemaSink.java 声明name = "jdbc-postgres"type = IOType.SINKconfigClass = JdbcSinkConfig.class
  • MariadbJdbcAutoSchemaSink.java 声明name = "jdbc-mariadb"
  • ClickHouse 与 SQLite 的实现类(ClickHouseJdbcAutoSchemaSink.java、SqliteJdbcAutoSchemaSink.java)遵循同样的模式。

各模块的 pom.xml 只额外引入对应数据库的 JDBC 驱动,例如 clickhouse/pom.xml 以 runtime 作用域引入ru.yandex.clickhouse:clickhouse-jdbc,并依赖pulsar-io-jdbc-core。这意味着:连接器与目标数据库的适配逻辑完全通用,差异仅在于驱动与注册的 sink 名称

二、配置属性详解:8 个核心参数与源码对照

所有 JDBC sink 连接器共享同一套配置结构,由 JdbcSinkConfig.java 定义。配置类使用@FieldDoc注解标注每个字段的必填性、默认值与帮助信息,Pulsar 管理工具据此生成连接器元数据。

属性类型必填默认值说明
userNameString空字符串""连接jdbcUrl指定数据库所使用的用户名。注意:userName区分大小写。
passwordString空字符串""连接jdbcUrl指定数据库所使用的密码。注意:password区分大小写。
jdbcUrlString空字符串""连接器要连接的数据库 JDBC URL。
tableNameString空字符串""连接器写入消息的目标数据表名称。
nonKeyString空字符串""逗号分隔的字段列表,用于 UPDATE 事件中 SET 子句的字段。
keyString空字符串""逗号分隔的字段列表,用于 UPDATE 与 DELETE 事件 WHERE 条件的字段。
timeoutMsint500JDBC 操作超时时间(毫秒)。
batchSizeint200写入数据库的批量大小(一次批量提交的记录数)。

2.1 源码对照:参数如何被解析与校验

open()阶段(JdbcAbstractSink.java),连接器依次执行:

  1. JdbcSinkConfig.load(config)把传入的 Map(JSON/YAML 反序列化结果)映射为配置对象;
  2. 校验jdbcUrl非空,为空时抛出IllegalArgumentException("Required jdbc Url not set.")
  3. 通过 JdbcUtils.getDriverClassName() 依据 URL 前缀自动匹配驱动类,再Class.forName加载驱动,DriverManager.getConnection建立连接;
  4. 连接建立后立即setAutoCommit(false),事务提交交由 flush 逻辑统一控制。

驱动识别由 JdbcDriverType.java 枚举完成,采用"URL 前缀 -> 驱动类"的映射策略。虽然仓库中注册了大量驱动(MySQL、DB2、Oracle、SQL Server、H2 等),但当前发布形态下仅打包 ClickHouse、MariaDB、PostgreSQL、SQLite 四种,其它条目服务于测试或未来扩展。

2.2 key / nonKey 的语义:Update 与 Delete 的字段分工

keynonKey直接决定了自动生成的 UPDATE / DELETE SQL 形态(见 JdbcUtils.java):

  • 配置了nonKey,才会生成并预编译UPDATE <table> SET <nonKey列>=?, ... WHERE <key列>=?
  • 配置了key,才会生成并预编译DELETE FROM <table> WHERE <key列>=?
  • INSERT 语句始终生成:INSERT INTO <table>(<全列>) VALUES(?, ...)

未配置nonKey/key时,连接器仅执行 INSERT,这正与文档"目前支持 INSERT、DELETE 和 UPDATE 操作"的说明相呼应——三种操作的能力上限取决于用户是否在配置中声明字段分工。

三、四种数据库的配置示例

连接器配置文件既可以用 JSON 也可以写成 YAML,运行期均会被解析为Map<String, Object>交给JdbcSinkConfig.load(Map)

3.1 ClickHouse

JSON:

{ "configs": { "userName": "clickhouse", "password": "password", "jdbcUrl": "jdbc:clickhouse://localhost:8123/pulsar_clickhouse_jdbc_sink", "tableName": "pulsar_clickhouse_jdbc_sink" } }

YAML:

tenant: "public" namespace: "default" name: "jdbc-clickhouse-sink" topicName: "persistent://public/default/jdbc-clickhouse-topic" sinkType: "jdbc-clickhouse" configs: userName: "clickhouse" password: "password" jdbcUrl: "jdbc:clickhouse://localhost:8123/pulsar_clickhouse_jdbc_sink" tableName: "pulsar_clickhouse_jdbc_sink"

3.2 MariaDB

JSON:

{ "configs": { "userName": "mariadb", "password": "password", "jdbcUrl": "jdbc:mariadb://localhost:3306/pulsar_mariadb_jdbc_sink", "tableName": "pulsar_mariadb_jdbc_sink" } }

YAML:

tenant: "public" namespace: "default" name: "jdbc-mariadb-sink" topicName: "persistent://public/default/jdbc-mariadb-topic" sinkType: "jdbc-mariadb" configs: userName: "mariadb" password: "password" jdbcUrl: "jdbc:mariadb://localhost:3306/pulsar_mariadb_jdbc_sink" tableName: "pulsar_mariadb_jdbc_sink"

3.3 PostgreSQL

使用 JDBC PostgreSQL sink 之前,需先通过下述任一方式创建配置文件。

JSON:

{ "configs": { "userName": "postgres", "password": "password", "jdbcUrl": "jdbc:postgresql://localhost:5432/pulsar_postgres_jdbc_sink", "tableName": "pulsar_postgres_jdbc_sink" } }

YAML:

tenant: "public" namespace: "default" name: "jdbc-postgres-sink" topicName: "persistent://public/default/jdbc-postgres-topic" sinkType: "jdbc-postgres" configs: userName: "postgres" password: "password" jdbcUrl: "jdbc:postgresql://localhost:5432/pulsar_postgres_jdbc_sink" tableName: "pulsar_postgres_jdbc_sink"

关于如何端到端使用该连接器(搭建 PostgreSQL 集群、建表、上传 schema、创建 sink),完整操作步骤见 connect Pulsar to PostgreSQL。

3.4 SQLite

SQLite 是嵌入式数据库,通常无需用户名密码,因此示例配置最精简:

JSON:

{ "configs": { "jdbcUrl": "jdbc:sqlite:db.sqlite", "tableName": "pulsar_sqlite_jdbc_sink" } }

YAML:

tenant: "public" namespace: "default" name: "jdbc-sqlite-sink" topicName: "persistent://public/default/jdbc-sqlite-topic" sinkType: "jdbc-sqlite" configs: jdbcUrl: "jdbc:sqlite:db.sqlite" tableName: "pulsar_sqlite_jdbc_sink"

四、源码级原理:连接器内部是如何工作的

4.1 表结构与 SQL 的自动发现与构建

连接器在open()中通过JdbcUtils.getTableId(connection, tableName)DatabaseMetaData校验目标表是否存在(不存在直接抛异常),随后getTableDefinition(...)读取该表的全部列名、SQL 类型(java.sql.Types)与列位置,并按key/nonKey配置把列划分为 keyColumns 与 nonKeyColumns。基于这份表定义,buildInsertSql/buildUpdateSql/buildDeleteSql自动生成三类PreparedStatement并预编译。因此目标表必须预先创建,连接器不会自动建表。

4.2 消息字段到列的绑定

BaseJdbcAutoSchemaSink.java 负责把GenericRecord消息绑定到 PreparedStatement:

  • INSERT:绑定表的全部列;
  • DELETE:只绑定 keyColumns;
  • UPDATE:绑定 nonKeyColumns + keyColumns。

绑定过程中按值类型分发到setInt / setLong / setDouble / setFloat / setBoolean / setString / setShort,其余类型会抛出 "Not support value type" 异常;字段缺失(JSON schema 省略字段导致的 NPE)或值为 null 时调用setNull(index, sqlType)写入数据库 NULL。

4.3 ACTION 机制:如何触发 Insert / Update / Delete

连接器通过消息属性ACTION决定写入方式(JdbcAbstractSink.java):

  • 属性未设置或值为INSERT:执行插入;
  • 值为UPDATE:执行更新(依赖nonKey+key配置);
  • 值为DELETE:执行删除(依赖key配置);
  • 其它值:抛出IllegalArgumentException

4.4 批量写入与定时 flush

batchSizetimeoutMs共同构成写入节奏:

  • 每收到一条消息先加入incomingList,累计达到batchSize时立即调度一次 flush(scheduleAtFixedRate之外再schedule(..., 0ms));
  • 与此同时,open()中创建的调度线程以timeoutMs为周期定时执行 flush,保证低流量时数据也能及时落库;
  • flush 采用incomingList/swapList双缓冲与AtomicBoolean互斥,逐条执行对应 PreparedStatement 后统一connection.commit(),成功则Record::ack,任一失败则整批Record::fail

这种"定时 + 定量"双重触发机制,正是batchSize=200timeoutMs=500这两个默认值在低延迟与吞吐之间取得平衡的工程实现。

五、端到端实战参考:从 Topic 到 PostgreSQL 表

以 PostgreSQL 为例,site2/docs/io-quickstart.md 给出了完整的落地路径,核心步骤为:

  1. 准备数据库与表:通过 Docker 启动 PostgreSQL,执行create table if not exists pulsar_postgres_jdbc_sink (id serial PRIMARY KEY, name VARCHAR(255) NOT NULL, ...)
  2. 编写配置文件:创建pulsar-postgres-jdbc-sink.yaml并置于pulsar/connectors目录,内容即上文 3.3 节的configs段;
  3. 为 topic 上传 AVRO schema:用bin/pulsar-admin schemas upload pulsar-postgres-jdbc-sink-topic -f ./connectors/avro-schema,再用bin/pulsar-admin schemas get校验;
  4. 创建 sink
bin/pulsar-admin sinks create \ --archive ./connectors/pulsar-io-jdbc-postgres-<version>.nar \ --inputs pulsar-postgres-jdbc-sink-topic \ --name pulsar-postgres-jdbc-sink \ --sink-config-file ./connectors/pulsar-postgres-jdbc-sink.yaml \ --parallelism 1

命令执行后,Pulsar 会以 Pulsar Function 的形态运行该 sink,将pulsar-postgres-jdbc-sink-topic中生产的消息持续写入 PostgreSQL 表pulsar_postgres_jdbc_sink。其中--archive指定连接器 NAR 包路径,--inputs为输入 topic(可逗号分隔多个),--name为 sink 名称,--sink-config-file指向 YAML 配置,--parallelism指定并行实例数。

六、使用注意事项

  • 目标表必须预先存在:连接器通过DatabaseMetaData校验表,只读表结构并自动生成 SQL,不会自动建表、改表;
  • userName/password区分大小写:两个字段在配置中均按原样传入连接属性;
  • UPDATE / DELETE 需要显式配置字段:不配置nonKey/key时连接器只具备 INSERT 能力;
  • 支持的数据类型有限:绑定值仅覆盖整数、长整型、浮点、布尔、字符串与短整型,其余 Java 类型需要扩展BaseJdbcAutoSchemaSink
  • 批量失败语义:flush 以整批为单位提交与确认(ack),整批任一语句失败则整批标记失败(fail),不存在单条部分确认;
  • 事务与连接:连接全程autoCommit=false,关闭连接器时会先commit()再关闭连接并关停 flush 线程。

以上行为均可直接对照 JdbcAbstractSink.java、JdbcUtils.java 与 JdbcSinkConfig.java 验证,建议在接入新数据库或排查写入问题时优先查阅这三处源码。

  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询