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

Apache Pulsar Cassandra Sink 连接器:配置详解与从 Topic 到 Cassandra 的数据写入实战

发布时间:2026/9/23 14:56:12 来源:云帆数科 栏目:资讯中心
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:12

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

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

PermissionsDispatcher 的 Java 使用指南:注解驱动的 Android 运行时权限处理完整实践
PermissionsDispatcher 的 Java 使用指南:注解驱动的 Android 运行时权限处理完整实践

PermissionsDispatcher 的 Java 使用指南:注解驱动的 Android 运行时权限处理完整实践 【免费下载链接】PermissionsDispatcher A declarative API to handle Android runtime permissions. 项目地址: https://gitcode.com/gh_mirrors/pe/PermissionsDispatcher … · 2026/9/23 14:56:04

福大易班源码解析:3个坑点避开,后端代码直接跑通
福大易班源码解析:3个坑点避开,后端代码直接跑通

福大易班源码解析:3个坑点避开,后端代码直接跑通 刚接手福大易班这类校园社区项目的后端维护时,最崩溃的不是需求多,而是从网上复制来的代码片段,丢进本地环境就报错。明明照着教程写的,为什么别人能跑,你这里却满屏红字?别急,这通常不是你的锅,而… · 2026/9/23 15:35:22

脑电数据分析利器:EEGLAB从预处理到ERP/频谱/时频分析实战指南
脑电数据分析利器:EEGLAB从预处理到ERP/频谱/时频分析实战指南

EEGLAB我从研究生阶段一直用到现在,中间换过好几个数据处理工具,最后还是老老实实回到它上面。这个工具箱确实不是最漂亮的那个,上手也有点门槛,但它把脑电数据分析从头到尾的环节都串起来了,网上随时能搜到教程&#… · 2026/9/23 15:35:16

华为技术专家揭秘百万年薪构成与职场跃迁策略
华为技术专家揭秘百万年薪构成与职场跃迁策略

1. 薪资数字背后的职场密码那天收到银行短信提醒时,我盯着屏幕反复数了三遍小数点前的位数——1,002,415.13这个数字确实没看错。作为在华为体系内深耕七年的技术专家,这个薪资数字既是对过往付出的肯定,也折射出科技行业顶尖人才的薪酬现状。… · 2026/9/23 15:35:16

花边边框简单漂亮图片生成提速80%的最佳实践
花边边框简单漂亮图片生成提速80%的最佳实践

花边边框简单漂亮图片生成提速80%的最佳实践 官方文档翻了三遍还是不知道哪里卡脖子?别急,今天直接上干货。很多开发者在做 花边边框简单漂亮图片 时,都遇到过渲染慢、内存爆的问题。其实核心就在于纹理加载和绘制批处理的细节。这篇不讲虚的,只讲… · 2026/9/23 15:35:10

XP系统关机后自动重启排查指南:从软件到硬件全流程
XP系统关机后自动重启排查指南:从软件到硬件全流程

简介:这份PDF文档专门解决Windows XP系统无法正常关机、关机后自动重启的经典问题,面向电脑维修人员、企业IT运维及仍在维护老机器的技术爱好者。资源共1个PDF文件,压缩包仅19KB,便携易用。文档首先解释Windows关机过程要完成的写… · 2026/9/23 15:35:10

双活数据中心端到端架构全解析:从存储到数据库的容灾设计
双活数据中心端到端架构全解析:从存储到数据库的容灾设计

简介:双活数据中心解决方案.pptx 是一份面向灾备架构师、运维工程师及企业IT决策者的技术讲解资料,聚焦两地三中心场景下的业务连续性与数据零丢失设计。基于华为双活数据中心端到端技术架构,资源从存储、应用、网络三个层面展开:… · 2026/9/23 15:35:10

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

了解更多?预约专属演示

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

企业微信二维码