首页/新闻资讯/正文详情

Apache Pulsar JDBC Sink Connector 完全指南:将 Topic 消息持久化到 ClickHouse / MariaDB / PostgreSQL / SQLite

发布时间:2026/9/23 10:00:04 来源:云帆数科 栏目:资讯中心
Apache Pulsar JDBC Sink Connector 完全指南:将 Topic 消息持久化到 ClickHouse / MariaDB / PostgreSQL / SQLite
消息队列后端流处理【免费下载链接】pulsarApache 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 类型sinkTypeClickHousepulsar-io/jdbc/clickhousejdbc-clickhouseMariaDBpulsar-io/jdbc/mariadbjdbc-mariadbPostgreSQLpulsar-io/jdbc/postgresjdbc-postgresSQLitepulsar-io/jdbc/sqlitejdbc-sqlite从源码看四个专属模块的实现类本身几乎是空壳——它们只通过Connector注解声明连接器名称、类型与配置类实际逻辑全部继承自核心模块。例如PostgresJdbcAutoSchemaSink.java 声明name jdbc-postgres、type IOType.SINK、configClass JdbcSinkConfig.classMariadbJdbcAutoSchemaSink.java 声明name jdbc-mariadbClickHouse 与 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 条件的字段。timeoutMsint否500JDBC 操作超时时间毫秒。batchSizeint否200写入数据库的批量大小一次批量提交的记录数。2.1 源码对照参数如何被解析与校验在open()阶段JdbcAbstractSink.java连接器依次执行JdbcSinkConfig.load(config)把传入的 MapJSON/YAML 反序列化结果映射为配置对象校验jdbcUrl非空为空时抛出IllegalArgumentException(Required jdbc Url not set.)通过 JdbcUtils.getDriverClassName() 依据 URL 前缀自动匹配驱动类再Class.forName加载驱动DriverManager.getConnection建立连接连接建立后立即setAutoCommit(false)事务提交交由 flush 逻辑统一控制。驱动识别由 JdbcDriverType.java 枚举完成采用URL 前缀 - 驱动类的映射策略。虽然仓库中注册了大量驱动MySQL、DB2、Oracle、SQL Server、H2 等但当前发布形态下仅打包 ClickHouse、MariaDB、PostgreSQL、SQLite 四种其它条目服务于测试或未来扩展。2.2 key / nonKey 的语义Update 与 Delete 的字段分工key与nonKey直接决定了自动生成的 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运行期均会被解析为MapString, Object交给JdbcSinkConfig.load(Map)。3.1 ClickHouseJSON{ configs: { userName: clickhouse, password: password, jdbcUrl: jdbc:clickhouse://localhost:8123/pulsar_clickhouse_jdbc_sink, tableName: pulsar_clickhouse_jdbc_sink } }YAMLtenant: 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_sink3.2 MariaDBJSON{ configs: { userName: mariadb, password: password, jdbcUrl: jdbc:mariadb://localhost:3306/pulsar_mariadb_jdbc_sink, tableName: pulsar_mariadb_jdbc_sink } }YAMLtenant: 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_sink3.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 } }YAMLtenant: 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 SQLiteSQLite 是嵌入式数据库通常无需用户名密码因此示例配置最精简JSON{ configs: { jdbcUrl: jdbc:sqlite:db.sqlite, tableName: pulsar_sqlite_jdbc_sink } }YAMLtenant: 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消息绑定到 PreparedStatementINSERT绑定表的全部列DELETE只绑定 keyColumnsUPDATE绑定 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执行更新依赖nonKeykey配置值为DELETE执行删除依赖key配置其它值抛出IllegalArgumentException。4.4 批量写入与定时 flushbatchSize与timeoutMs共同构成写入节奏每收到一条消息先加入incomingList累计达到batchSize时立即调度一次 flushscheduleAtFixedRate之外再schedule(..., 0ms)与此同时open()中创建的调度线程以timeoutMs为周期定时执行 flush保证低流量时数据也能及时落库flush 采用incomingList/swapList双缓冲与AtomicBoolean互斥逐条执行对应 PreparedStatement 后统一connection.commit()成功则Record::ack任一失败则整批Record::fail。这种定时 定量双重触发机制正是batchSize200、timeoutMs500这两个默认值在低延迟与吞吐之间取得平衡的工程实现。五、端到端实战参考从 Topic 到 PostgreSQL 表以 PostgreSQL 为例site2/docs/io-quickstart.md 给出了完整的落地路径核心步骤为准备数据库与表通过 Docker 启动 PostgreSQL执行create table if not exists pulsar_postgres_jdbc_sink (id serial PRIMARY KEY, name VARCHAR(255) NOT NULL, ...)编写配置文件创建pulsar-postgres-jdbc-sink.yaml并置于pulsar/connectors目录内容即上文 3.3 节的configs段为 topic 上传 AVRO schema用bin/pulsar-admin schemas upload pulsar-postgres-jdbc-sink-topic -f ./connectors/avro-schema再用bin/pulsar-admin schemas get校验创建 sinkbin/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不存在单条部分确认事务与连接连接全程autoCommitfalse关闭连接器时会先commit()再关闭连接并关停 flush 线程。以上行为均可直接对照 JdbcAbstractSink.java、JdbcUtils.java 与 JdbcSinkConfig.java 验证建议在接入新数据库或排查写入问题时优先查阅这三处源码。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar Solr Sink Connector 配置与源码剖析将 Topic 消息持久化到 Solr CollectionApache Pulsar Solr Sink Connector 配置与源码剖析将 Topic 消息持久化到 Solr Collection Solr si消息队列后端流处理Apache Pulsar Redis Sink Connector 完全指南将 Topic 消息实时写入 RedisApache Pulsar Redis Sink Connector 完全指南将 Topic 消息实时写入 Redis 本篇技术指南以 Apache Puls消息队列后端流处理Apache Pulsar HDFS2 Sink Connector 完全指南从 Pulsar Topic 落盘 HDFS 的配置、部署与源码剖析Apache Pulsar HDFS2 Sink Connector 完全指南从 Pulsar Topic 落盘 HDFS 的配置、部署与源码剖析 HDFS2消息队列后端流处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关推荐

电子元器件控制信号分类与RS485自收发电路可靠性分析
电子元器件控制信号分类与RS485自收发电路可靠性分析

设备做多了以后,你会发现一个很有意思的现象:一张原理图拿过来,甭管上面密密麻麻摆了多少个芯片、多少颗电阻电容,真正决定这块板子“性格”的,往往就是那几根控制信号线。电子元器件怎么工作、什么时候工作、以什么方… · 2026/9/23 9:59:58

3步搞定怎样无线连接打印机:手写实现避坑指南
3步搞定怎样无线连接打印机:手写实现避坑指南

3步搞定怎样无线连接打印机:手写实现避坑指南 看了一堆教程还是不会写项目?别慌,这通常是把“操作指南”当成了“编程逻辑”。很多开发者在面对【怎样无线连接打印机】这种偏硬件交互的软技能时,容易陷入“点鼠标”的思维定式,忽略了底层协议与代码实现… · 2026/9/23 9:59:51

SHIO电子管箱头:从电路设计到调试实践的英式音色之旅
SHIO电子管箱头:从电路设计到调试实践的英式音色之旅

Tubes&Tone这个名字出现在大家的订阅列表里已经有一阵子了。这些年我们做的事情,说简单点就是跟电子管和音色打交道——修过几十台老箱子,改过数不清的效果器,也拆过不少让人心疼的贵重设备。但直到SHIO装进箱壳、从喇叭里发出第一声干净… · 2026/9/23 9:59:51

2024 CSP-J初赛真题深度解析:考点拆解与复习路径
2024 CSP-J初赛真题深度解析:考点拆解与复习路径

简介:这份资源是2024年信息学奥赛CSP-J初赛真题的详细分析文档,面向备战CSP-J的中学生、少儿编程学习者及竞赛指导教师,帮助读者系统梳理初赛考点、理解命题思路并查漏补缺。压缩包内共1个docx文件,约121KB,内容以文字… · 2026/9/23 10:55:47

卡方检验表避坑指南:3个实战案例教你搞定API变更
卡方检验表避坑指南:3个实战案例教你搞定API变更

卡方检验表避坑指南:3个实战案例教你搞定API变更 版本升级后 API 全变了?别慌,这篇避坑指南专治各种统计库升级导致的“水土不服”。 在水利工程的数据分析里,卡方检验表是检验独立性、拟合优度的核心工具。很多老手都遇到过:项目从… · 2026/9/23 10:55:47

长风破浪会有时高频面试题全解:应届生避坑指南
长风破浪会有时高频面试题全解:应届生避坑指南

长风破浪会有时高频面试题全解:应届生避坑指南 报错堆栈长得像天书,面试官一句“长风破浪会有时”让你解释底层逻辑,你脑子瞬间空白?别慌。… · 2026/9/23 10:55:47

中文短文本情感分析实战:基于PyTorch的LSTM模型搭建与调优指南
中文短文本情感分析实战:基于PyTorch的LSTM模型搭建与调优指南

简介:一套基于LSTM的中文短文本情感分析项目源码,主要面向需要完成期末大作业、课程设计的高校学生,也适合刚接触深度学习文本分类的Python开发者用于练手与拓展。压缩包共14个文件、大小约1.96MB,内部包含Python源码脚本、训练与… · 2026/9/23 10:55:40

基于SQLite与三层容错架构的内容获取工作流重构实践
基于SQLite与三层容错架构的内容获取工作流重构实践

1. 内容获取工作流的核心痛点与重构思路做过内容批量采集的人都有一个共识:真正让人头疼的从来不是“能不能下载”,而是“下载过程稳不稳”。douyin-downloader 这类工具在圈子里流传已久,早期版本大多走的是单链路请求——解析一个视频 ID&a… · 2026/9/23 10:55:40

数字绘画笔刷管理:从参数原理到高效工作流
数字绘画笔刷管理:从参数原理到高效工作流

1. 项目概述:一个手绘爱好者的工具进化史五年前刚接触数字绘画时,我和大多数新手一样陷入"笔刷收藏癖"的怪圈——硬盘里囤积了上百套笔刷却从未认真使用过任何一套。直到有次接商业项目时,甲方要求用特定风格的笔触完成整套插画&am… · 2026/9/23 10:55:40

3招搞定手机怎么下载微信面试难题实战项目解析
3招搞定手机怎么下载微信面试难题实战项目解析

3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A… · 2026/9/23 0:00:03

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型
你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 面试被问“高并发下如何保证消息不丢失”,你张口就是“用Redis”,结果面试官追问“如果Redis宕机了怎么办”,你瞬间卡壳。这种场景太常见了,很多新手在背八股文时,只记住了技术名词… · 2026/9/23 0:00:29

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧
Win7无线热点配置工具源码解析:解决API失效的3个实战技巧

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧 Win7无线热点配置工具在Win10/11上跑不动?不是你的问题,是版本升级后 API 全变了。很多老项目里的 netsh wlan… · 2026/9/23 0:00:36

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码