ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

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

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),仅供参考
返回列表