Apache Pulsar Cassandra Sink 连接器:配置详解与从 Topic 到 Cassandra 的数据写入实战
Apache Pulsar Cassandra Sink 连接器配置详解与从 Topic 到 Cassandra 的数据写入实战【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsarApache Pulsar 的 Cassandra sink connector 负责把 Pulsar topic 中的消息拉取出来并写入到 Cassandra 集群的表中是典型的消息落库场景组件。本篇指南以仓库中的 io-cassandra-sink.md 为主体结合pulsar-io/cassandra模块源码、集成测试与 io-quickstart.md 完整流程系统讲解其全部配置项、JSON/YAML 配置写法、从启动 Cassandra 到创建/校验/删除 sink 的完整命令链并深入到CassandraAbstractSink源码揭示其异步写入与消息确认机制。读完本文你将能独立完成Pulsar → Cassandra管道的搭建与排障。一、Cassandra sink 是什么Cassandra sink connector 是 Pulsar IO 框架中的内置连接器之一SinkType.CASSANDRA它的职责非常单一且明确读取 Pulsar 输入 topic 中的每条消息将其 key/value 写入 Cassandra 的指定列族表。从源码结构看该连接器由pulsar-io/cassandra模块提供包含三个核心类CassandraSinkConfig.java连接器配置模型负责从 YAML 文件或 Map 中加载五个必填参数CassandraAbstractSink.java抽象的 Sink 实现封装了连接建立、INSERT语句预编译、异步写入与 ack/fail 回调CassandraStringSink.java默认实现类将消息按字符串处理写入相同的 key/value 对。其中CassandraStringSink通过Connector(name cassandra, type IOType.SINK, configClass CassandraSinkConfig.class)注解注册name cassandra正是pulsar-admin sinks create --sink-type cassandra中使用的类型名。模块由nifi-nar-maven-plugin打包为 NAR 归档见 pulsar-io/cassandra/pom.xml既可作为内置连接器使用也可作为独立 NAR 分发。二、配置属性详解Cassandra sink 的全部配置集中在一个配置文件里共5 个必填属性。下表完整列出属性、类型、是否必填、默认值与说明名称类型必填默认值说明rootsString是空字符串要连接的 Cassandra 主机列表多个主机用逗号分隔格式为host:port。keyspaceString是空字符串用于写入 Pulsar 消息的 keyspace。注意keyspace必须在启动 Cassandra sink 之前预先创建。keynameString是空字符串Cassandra 列族中用于存储 Pulsar 消息 key 的列名。如果 Pulsar 消息没有关联 key则使用消息 value 作为 key 写入该列。columnFamilyString是空字符串Cassandra 列族表名称。注意columnFamily必须在启动 Cassandra sink 之前预先创建。columnNameString是空字符串Cassandra 列族中用于存储 Pulsar 消息 value 的列名。源码级校验逻辑这些必填约束在运行时由 CassandraSinkConfig.java 的FieldDoc(required true, defaultValue )注解声明并在 CassandraAbstractSink.java 的open()方法中强制执行if (cassandraSinkConfig.getRoots() null || cassandraSinkConfig.getKeyspace() null || cassandraSinkConfig.getKeyname() null || cassandraSinkConfig.getColumnFamily() null || cassandraSinkConfig.getColumnName() null) { throw new IllegalArgumentException(Required property not set.); }也就是说五个属性缺一不可漏配任何一个sink 在启动阶段就会直接抛出IllegalArgumentException。roots的解析逻辑也值得注意createClient()先把roots按逗号切分为主机列表再把每个host:port拆开逐个addContactPoint()只有当该项包含端口hostPort.length 1时才调用withPort()若未写端口则使用 DataStax Java Driver 的默认端口 9042。两个必须预创建的前提keyspace与columnFamily是 sink 无法自行创建的open()中执行的是session.execute(USE keyspace)切换 keyspace然后session.prepare(INSERT INTO columnFamily ( keyname , columnName ) VALUES (?, ?))预编译插入语句。如果 keyspace 或表不存在这两个调用都会失败。因此正确的使用顺序是先在 Cassandra 侧建好 keyspace 和表再启动 sink。三、配置文件示例JSON 与 YAML 两种写法原文档给出两种配置文件写法均完整继承如下。JSON 格式{ configs: { roots: localhost:9042, keyspace: pulsar_test_keyspace, columnFamily: pulsar_test_table, keyname: key, columnName: col } }YAML 格式configs: roots: localhost:9042 keyspace: pulsar_test_keyspace columnFamily: pulsar_test_table keyname: key columnName: col补充说明两点在 io-quickstart.md 的示例中JSON 也出现过不带configs外层包装的扁平写法即roots/keyspace等直接平铺在顶层。这两种形式在仓库文档中均被使用底层配置加载时都会被反序列化为CassandraSinkConfig对应的 MapCassandraSinkConfig.load(MapString, Object map)通过 Jackson 完成绑定。配置文件中五个键名必须与属性名完全一致大小写敏感Jackson 按字段名直接映射。四、前置准备启动 Cassandra 并创建 keyspace 与表在创建 sink 之前需要先有一个可用的 Cassandra 集群以及上面反复强调的 keyspace 与表。仓库的 io-quickstart.md 提供了基于 Docker 的完整步骤1. 启动单节点 Cassandra 集群docker run -d --rm --namecassandra -p 9042:9042 cassandra2. 确认进程与集群状态docker ps docker logs cassandra docker exec cassandra nodetool statusnodetool status正常时输出类似Datacenter: datacenter1 StatusUp/Down |/ StateNormal/Leaving/Joining/Moving -- Address Load Tokens Owns (effective) Host ID Rack UN 172.17.0.2 103.67 KiB 256 100.0% af0e4b2f-84e0-4f0b-bb14-bd5f9070ff26 rack1UN表示节点处于 Up在线且 Normal正常状态。3. 进入 cqlsh 创建 keyspace 与表$ docker exec -ti cassandra cqlsh localhost Connected to Test Cluster at localhost:9042. [cqlsh 5.0.1 | Cassandra 3.11.2 | CQL spec 3.4.4 | Native protocol v4] Use HELP for help. cqlshcqlsh CREATE KEYSPACE pulsar_test_keyspace WITH replication {class:SimpleStrategy, replication_factor:1}; cqlsh USE pulsar_test_keyspace; cqlsh:pulsar_test_keyspace CREATE TABLE pulsar_test_table (key text PRIMARY KEY, col text);这里创建的pulsar_test_keyspace与pulsar_test_table必须与配置文件中keyspace、columnFamily两个字段保持完全一致。表结构要求也很简单一张含两个文本列的表其中 key 列是主键分别对应配置里的keyname与columnName。仓库的集成测试 CassandraSinkTester.java 正是用同样的 CQLSimpleStrategyreplication_factor:1key text PRIMARY KEY, col text初始化测试环境可以作为最小可用表结构的权威参考。五、创建并运行 Cassandra sink准备好配置文件例如examples/cassandra-sink.yml内容即第三节中的 YAML后使用 Pulsar 的 Connector Admin CLI 创建 sink。以下命令将创建一个名为cassandra-test-sink、类型为cassandra、消费 topictest_cassandra的 sinkbin/pulsar-admin sinks create \ --tenant public \ --namespace default \ --name cassandra-test-sink \ --sink-type cassandra \ --sink-config-file examples/cassandra-sink.yml \ --inputs test_cassandra各参数含义--tenant/--namespacesink 归属的租户与命名空间--namesink 实例名后续 get/status/delete 都靠它定位--sink-type连接器类型内置连接器的类型名取自源码Connector注解中的nameCassandra 即cassandra--sink-config-file第三节中准备的配置文件路径--inputssink 的输入 topic。命令执行成功后Pulsar 会以 Pulsar Function 的形式运行该 sink它从 topictest_cassandra消费消息并写入 Cassandra 表pulsar_test_table。此命令需要在一个启用了 Functions Worker 的 Pulsar 集群上执行如 standalone 模式。检查 sink 信息与状态获取 sink 信息bin/pulsar-admin sinks get \ --tenant public \ --namespace default \ --name cassandra-test-sink输出示例注意其中的className正是源码中的CassandraStringSinkconfigs即我们配置的五个参数{ tenant: public, namespace: default, name: cassandra-test-sink, className: org.apache.pulsar.io.cassandra.CassandraStringSink, inputSpecs: { test_cassandra: { isRegexPattern: false } }, configs: { roots: localhost:9042, keyspace: pulsar_test_keyspace, columnFamily: pulsar_test_table, keyname: key, columnName: col }, parallelism: 1, processingGuarantees: ATLEAST_ONCE, retainOrdering: false, autoAck: true, archive: builtin://cassandra }检查运行状态bin/pulsar-admin sinks status \ --tenant public \ --namespace default \ --name cassandra-test-sink状态输出中值得关注三个计数器numReadFromPulsar已从 Pulsar 读取的消息数numWrittenToSink已成功写入 Cassandra 的消息数numSystemExceptions/numSinkExceptions系统异常与 sink 异常计数排障时首先查看。六、数据验证从生产消息到 cqlsh 查询1. 向输入 topic 生产 10 条消息for i in {0..9}; do bin/pulsar-client produce -m key-$i -n 1 test_cassandra; done2. 再次查看 sink 状态确认消息已被消费并写入bin/pulsar-admin sinks status \ --tenant public \ --namespace default \ --name cassandra-test-sink正常情况下numReadFromPulsar与numWrittenToSink都变为10{ numInstances : 1, numRunning : 1, instances : [ { instanceId : 0, status : { running : true, error : , numRestarts : 0, numReadFromPulsar : 10, numSystemExceptions : 0, latestSystemExceptions : [ ], numSinkExceptions : 0, latestSinkExceptions : [ ], numWrittenToSink : 10, lastReceivedTime : 1551685489136, workerId : c-standalone-fw-localhost-8080 } } ] }3. 回到 Cassandra 侧查表验证落库结果docker exec -ti cassandra cqlsh localhostcqlsh use pulsar_test_keyspace; cqlsh:pulsar_test_keyspace select * from pulsar_test_table;输出显示 10 条消息的 key 与 value 已一一对应写入本例中每条消息同时作为 key 与 value 落库key | col ---------------- key-5 | key-5 key-0 | key-0 key-9 | key-9 key-2 | key-2 key-1 | key-1 key-3 | key-3 key-6 | key-6 key-7 | key-7 key-4 | key-4 key-8 | key-8七、写入语义与消息确认机制源码解析为什么10 条消息写入后key 和 col 是相同的值答案在 CassandraStringSink.java 的extractKeyValue实现中Override public KeyValueString, String extractKeyValue(Recordbyte[] record) { String key record.getKey().orElseGet(() - new String(record.getValue())); return new KeyValue(key, new String(record.getValue())); }若消息带有 key则keyname列写入消息 key、columnName列写入消息 value若消息没有 key如本示例用pulsar-client produce直接发送的纯文本消息则退化为orElseGet分支用消息 value 本身作为 key—— 这正是原文档中如果 Pulsar 消息没有关联 key则使用消息 value 作为 key注释对应的实现。而整个写入流程定义在 CassandraAbstractSink.java 的write()方法中Override public void write(Recordbyte[] record) { KeyValueK, V keyValue extractKeyValue(record); BoundStatement bound statement.bind(keyValue.getKey(), keyValue.getValue()); ResultSetFuture future session.executeAsync(bound); Futures.addCallback(future, new FutureCallbackResultSet() { Override public void onSuccess(ResultSet result) { record.ack(); } Override public void onFailure(Throwable t) { record.fail(); } }, MoreExecutors.directExecutor()); }这里有三点值得深入理解预编译语句INSERT INTO columnFamily (keyname, columnName) VALUES (?, ?)在open()阶段一次性prepare之后每条消息只做bind与异步执行避免反复解析 CQL异步写入使用 DataStax Java Driver 的executeAsync异步写入并通过 Guava 的FutureCallback处理结果——写入成功回调record.ack()确认消息失败回调record.fail()触发重试/投递这是 Pulsar IO 框架标准的消息确认协议at-least-once 语义由于是先写 Cassandra、成功后才 ack若在写入成功但 ack 之前发生故障消息会被重新投递因此该 sink 表现为ATLEAST_ONCE与第五节get输出中的processingGuarantees一致。这意味着同一消息可能被重复写入如果你的下游查询依赖唯一性建议在表设计上利用key主键做幂等同 key 覆盖写。仓库的集成测试 CassandraSinkTester.java 对这一行为做了端到端校验向 sink 写入若干 key/value 后通过SELECT * FROM table拉回全表断言kvs.size()与行数相等、且每行的 value 与期望值一致。八、清理删除 Cassandra sink不再需要该管道时执行bin/pulsar-admin sinks delete \ --tenant public \ --namespace default \ --name cassandra-test-sink注意sinks delete只删除 Pulsar 侧的 sink 实例不会删除 Cassandra 中已写入的数据落库数据如需清理需在 Cassandra 侧单独处理。九、参考与延伸阅读本文主体文档io-cassandra-sink.md完整端到端实战流程io-quickstart.md连接器实现源码CassandraAbstractSink.java、CassandraStringSink.java、CassandraSinkConfig.java模块构建配置pulsar-io/cassandra/pom.xml集成测试验证CassandraSinkTester.java【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsar创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关新闻

RTL8370N实战:8端口L2管理型交换芯片的硬件与配置指南

RTL8370N实战:8端口L2管理型交换芯片的硬件与配置指南

简介:RTL8370NI-VB-CG数据手册是一份面向交换机软硬件工程师及产品选型人员的芯片参考文档,围绕瑞昱8端口10/100/1000M自适应二层管理型交换控制器展开。手册逐项说明芯片的端口自动协商、全双工与半双工模式、802.1Q虚拟局域网划分、基于端口或数据流的…

2026/9/23 14:56:19 阅读更多 →
5分钟搞懂送流量活动:从语法到项目的速查手册

5分钟搞懂送流量活动:从语法到项目的速查手册

5分钟搞懂送流量活动:从语法到项目的速查手册 刚学完 Python 或 Java 的 if-else,是不是觉得脑子清醒得很?一上手要搭个“送流量活动”页面,立马卡壳。很多人卡在“我会写代码,但不知道怎么把它变成产品”这一步。…

2026/9/23 14:56:19 阅读更多 →
Red Arrow:基于 GObject Introspection 的 Apache Arrow Ruby 绑定安装与实战指南

Red Arrow:基于 GObject Introspection 的 Apache Arrow Ruby 绑定安装与实战指南

Red Arrow:基于 GObject Introspection 的 Apache Arrow Ruby 绑定安装与实战指南 【免费下载链接】arrow Apache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing 项目地址: https://gitcode.com/gh_mirrors/arro…

2026/9/23 14:56:19 阅读更多 →

最新新闻

菱形虚拟继承的原理

菱形虚拟继承的原理

目录 摘要: 一 :菱形继承的概念及问题 1:概念 2:问题 二:虚拟菱形继承 1:语法 2:原理 ①:菱形继承的内存分布 ②:虚拟菱形继承的内存分布 ③:偏移量…

2026/9/23 15:44:20 阅读更多 →
学术写作AI:破解黑话,提升论文可读性与影响力

学术写作AI:破解黑话,提升论文可读性与影响力

1. 项目概述:当学术写作遇上"人话革命"去年审阅某核心期刊投稿时,我遇到一篇让我哭笑不得的论文——作者用"基于多维度认知框架的跨模态表征重构"来描述"用不同方法分析数据",通篇充斥着"后现代性话语解构…

2026/9/23 15:44:20 阅读更多 →
LPDDR5内存训练全流程解析:从ZQ校准到周期重训练的工程实践

LPDDR5内存训练全流程解析:从ZQ校准到周期重训练的工程实践

简介:面向内存控制器设计与嵌入式系统开发工程师,系统讲解LPDDR5内存的初始化与完整训练流程。内容涵盖上电初始化时序、ZQ校准(含输出驱动器阻抗校准与CA/DQ ODT阻抗校准)、命令总线训练、WCK与CK对齐、WCK占空比训练、读门控训练…

2026/9/23 15:44:20 阅读更多 →
3个避坑技巧搞定人体器官分布图代码面试必问

3个避坑技巧搞定人体器官分布图代码面试必问

3个避坑技巧搞定人体器官分布图代码面试必问 复制来的代码跑不通,控制台一堆红字报错,这时候你是不是只想把电脑砸了?这种“看似能跑实则崩盘”的情况,在技术面试中简直是重灾区。很多候选人拿着网上抄的 SVG 或 Canvas…

2026/9/23 15:44:20 阅读更多 →
搞定空间寄语:前端高薪必备的5个高频面试题

搞定空间寄语:前端高薪必备的5个高频面试题

搞定空间寄语:前端高薪必备的5个高频面试题 别再用“Hello World”糊弄自己了。很多学员学完语法,对着空白文档发呆,根本不知道怎么把零散的代码拼成一个能跑的项目。更扎心的是,面试官问起 高频面试题…

2026/9/23 15:44:20 阅读更多 →
RBAC权限系统设计与认证授权实践指南

RBAC权限系统设计与认证授权实践指南

1. 认证授权基础概念解析认证(Authentication)和授权(Authorization)是每个后端开发者必须掌握的核心安全机制。认证解决"你是谁"的问题,就像进入公司大楼时需要刷工牌确认身份;授权则解决"…

2026/9/23 15:43:19 阅读更多 →

日新闻

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

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

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

2026/9/23 0:00:23 阅读更多 →
2k显示屏性能优化踩坑:版本升级后API全变了,这份源码解析救了我

2k显示屏性能优化踩坑:版本升级后API全变了,这份源码解析救了我

2k显示屏性能优化踩坑:版本升级后API全变了,这份源码解析救了我 刚把开发环境的显示器从1080P换到2K,跑老项目直接报错,版本升级后 API…

2026/9/23 0:01:25 阅读更多 →
3步搞定美眉图实战项目,告别官方文档抓不住重点

3步搞定美眉图实战项目,告别官方文档抓不住重点

3步搞定美眉图实战项目,告别官方文档抓不住重点 官方文档翻了三遍还是云里雾里?别急,美眉图在实战项目中常被用来做数据可视化,但它的原理比你想的简单。今天咱们直接上手,用一个完整的小项目把美眉图跑通,不再死磕那些冗长的理论说明。…

2026/9/23 0:01:25 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/23 4:55:02 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/23 4:49:06 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/23 9:53:41 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/23 9:53:40 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/23 9:53:40 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/23 9:53:40 阅读更多 →