尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

RisingWave 实时流式写入 Cassandra / ScyllaDB 完整实战指南

发布时间:2026/9/25 3:21:50

资讯中心
01
ARTICLE

RisingWave 实时流式写入 Cassandra / ScyllaDB 完整实战指南

RisingWave 实时流式写入 Cassandra / ScyllaDB 完整实战指南
数据库流处理后端数据工程【免费下载链接】risingwaveEvent streaming platform for agentic AI. Continuously ingest, transform, and serve event streams in real time, at scale.项目地址https://gitcode.com/gh_mirrors/ri/risingwave点击查看免费下载本指南以 RisingWave 仓库中 integration_tests/cassandra-and-scylladb-sink 目录的官方 Demo 为骨架讲解如何让 RisingWave 将物化视图Materialized View中的实时数据持续写入 Apache Cassandra 与 ScyllaDB涵盖环境搭建、建表、建 Sink、数据校验全流程。读完本文你将掌握 RisingWave Cassandra Sink 的全部配置参数、类型映射规则、容器化联调方法并能用仓库内现成的脚本与测试用例自行复现与验证。一、Demo 概览与工作原理该 Demo 展示的核心链路是RisingWave 通过内置的cassandraconnector把流式计算结果以追加写入append-only方式同步到 Cassandra 与 ScyllaDB。集群由docker compose一键拉起包含以下组件RisingWave 单机集群及其依赖PostgreSQL、MinIO、Grafana、Prometheus、消息队列复用 docker/docker-compose.yml 中定义的服务一个 datagen 连接器持续生成用户行为模拟数据一台 Apache Cassandra 4.0端口 9042与一台 ScyllaDB 5.1端口 9041内部仍为 9042作为 Sink 目标库。从源码结构看RisingWave 的 Sink 层通过统一的连接器框架将cassandra映射为CassandraSink实现见 src/connector/src/sink/mod.rs 与 src/connector/src/sink/remote.rs 中的{ Cassandra, CassandraSink, cassandra, [ cassandra.url ] }注册项。因此在同一份 SQL 中只需修改cassandra.url指向不同主机即可将同一份数据同时写入 Cassandra 与 ScyllaDB——二者都兼容 CQL 协议这也是本 Demo 能够一鱼两吃的关键。二、环境搭建一键启动集群进入 demo 目录并启动全部服务cd integration_tests/cassandra-and-scylladb-sink docker-compose up -dintegration_tests/cassandra-and-scylladb-sink/docker-compose.yml 中两个数据库容器的关键配置如下cassandra: image: cassandra:4.0 ports: - 9042:9042 environment: - CASSANDRA_CLUSTER_NAMEcloudinfra volumes: - ./prepare_cassandra_and_scylladb.sql:/prepare_cassandra_and_scylladb.sql scylladb: image: scylladb/scylla:5.1 ports: - 9041:9042 # 宿主 9041 已被 cassandra 占用故映射到 9041 environment: - CASSANDRA_CLUSTER_NAMEcloudinfra值得注意的是ScyllaDB 容器内部仍然监听 9042 端口CQL 默认端口但宿主机 9041 已被 Cassandra 占用因此映射到宿主 9041。RisingWave 侧的cassandra.url使用的是Docker 网络内部服务名cassandra:9042与scylladb:9042而不是宿主端口。三、在 Cassandra/ScyllaDB 侧准备 Keyspace 与表3.1 通过 cqlsh 手工建表README 标准流程分别登录两个数据库的 cqlsh# 进入 Cassandra docker compose exec cassandra cqlsh # 进入 ScyllaDB docker compose exec scylladb cqlsh依次执行建库建表语句CREATE KEYSPACE demo WITH replication {class: SimpleStrategy, replication_factor: 1}; use demo; CREATE table demo_bhv_table( user_id int primary key, target_id text, event_timestamp timestamp, );SimpleStrategyreplication_factor: 1适用于单节点演示环境生产环境建议按集群拓扑选用NetworkTopologyStrategy并设置合理副本数。3.2 仓库内置的一键初始化脚本仓库同时提供了自动化脚本 integration_tests/cassandra-and-scylladb-sink/prepare.sh等待 30 秒让数据库完成启动后依次对两个容器执行docker compose exec cassandra cqlsh -f prepare_cassandra_and_scylladb.sql docker compose exec scylladb cqlsh -f prepare_cassandra_and_scylladb.sql其中 integration_tests/cassandra-and-scylladb-sink/prepare_cassandra_and_scylladb.sql 除了demo_bhv_table外还创建了一张用于验证类型映射的cassandra_types表CREATE table cassandra_types ( types_id int primary key, c_boolean boolean, c_smallint smallint, c_integer int, c_bigint bigint, c_decimal decimal, c_real float, c_double_precision double, c_varchar text, c_bytea blob, c_date date, c_time time, c_timestamptz timestamp, c_interval duration );四、在 RisingWave 侧创建 Source 与物化视图4.1 创建 Sourcedatagen 持续模拟数据执行 integration_tests/cassandra-and-scylladb-sink/create_source.sql它创建两张表user_behaviors使用connector datagen内置连接器user_id按 11000 的序列递增其余字段随机生成datagen.rows.per.second 10控制每秒产生 10 行数据格式为FORMAT PLAIN ENCODE JSON。datagen 是 RisingWave 内置的模拟数据源无需外部系统即可持续产生流式数据非常适合联调 Sink 链路。cassandra_types一张普通表随后通过三条INSERT写入覆盖极值、边界值与特殊类型的行例如-9223372036854775807、9999-12-31、9990 year区间等用于验证各类 RisingWave 类型能否正确落到 Cassandra 对应类型。4.2 创建物化视图执行 integration_tests/cassandra-and-scylladb-sink/create_mv.sql从user_behaviors投影出三个字段CREATE MATERIALIZED VIEW bhv_mv AS SELECT user_id, target_id, event_timestamp FROM user_behaviors;物化视图会持续增量维护查询结果作为后续 Sink 的数据源——这正是 RisingWave“流上建仓、实时出数”的典型形态。五、创建 Sink一个连接器双写 Cassandra 与 ScyllaDB按顺序依次执行create_source.sql→create_mv.sql→create_sink.sql。核心的 integration_tests/cassandra-and-scylladb-sink/create_sink.sql 内容如下set sink_decouple false; CREATE SINK bhv_cassandra_sink FROM bhv_mv WITH ( connector cassandra, type append-only, force_append_onlytrue, cassandra.url cassandra:9042, cassandra.keyspace demo, cassandra.table demo_bhv_table, cassandra.datacenter datacenter1, ); CREATE SINK bhv_scylla_sink FROM bhv_mv WITH ( connector cassandra, type append-only, force_append_onlytrue, cassandra.url scylladb:9042, cassandra.keyspace demo, cassandra.table demo_bhv_table, cassandra.datacenter datacenter1, );5.1 参数逐项说明参数值含义connectorcassandra指定使用 Cassandra/ScyllaDB 连接器typeappend-only追加写入模式不做 upsert 语义force_append_onlytrue强制按 append-only 处理即使上游可能含更新也忽略其变更语义cassandra.urlcassandra:9042/scylladb:9042目标数据库地址Docker 网络内服务名:端口cassandra.keyspacedemo目标 Keyspacecassandra.tabledemo_bhv_table目标表名cassandra.datacenterdatacenter1Cassandra 驱动连接所用的数据中心名需与集群实际配置一致cassandra.url是连接器唯一必填属性见 src/connector/src/sink/remote.rs 中[ cassandra.url ]的注册声明。set sink_decouple false;表示关闭 Sink 解耦写入行为跟随事务提交执行便于 Demo 中即时校验。5.2 类型映射验证create_sink.sql后半段把cassandra_types表分别通过cassandra_types_sink和scylladb_types_sink两个 Sink 写入两个数据库的cassandra_types表用于端到端验证类型映射。RisingWave 侧类型与 Cassandra 侧的对应关系为boolean→booleansmallint→smallintinteger→intbigint→bigintdecimal→decimalreal→floatdouble precision→doublevarchar→textbytea→blobdate→datetime→timetimestamptz→timestampinterval→duration仓库在 e2e_test/sink/cassandra_sink.slt 中提供了等价的自动化回归用例CI 通过 ci/scripts/e2e-cassandra-sink-test.sh 驱动其中还覆盖了带引号的大小写敏感表名Test_uppercase的写入场景可作为生产环境核对字段映射的参考。六、校验写入结果6.1 手工查询验证等 datagen 持续灌入数据后重新进入 cqlsh 执行聚合查询select user_id, count(*) from demo.demo_bhv_table group by user_id;由于 datagen 以user_id为主键PRIMARY KEY(user_id)且每秒生成 10 行Cassandra 侧会按主键覆盖更新同一user_id的target_id与event_timestamp因此预期每个user_id对应一条记录共 1000 个user_id。6.2 脚本化自动校验仓库提供了 integration_tests/cassandra-and-scylladb-sink/sink_check.py对demo.demo_bhv_table与demo.cassandra_types两张表、两个数据库逐一执行select count(*)并通过assert rows 1判定写入成功任一表为空即报错退出。运行方式python3 sink_check.py该脚本以docker compose exec db cqlsh -e sql的方式封装校验逻辑任何失败案例会汇总打印Data check failed for case ...并以非零码退出可直接接入 CI 门禁。七、清理与注意事项端口冲突Cassandra 与 ScyllaDB 都监听 9042docker-compose 中 ScyllaDB 已映射到宿主 9041勿再为两个容器分配相同宿主端口。启动时序两个数据库首次启动需要约 30 秒初始化prepare.sh与sink_check.py都内置了sleep(30)等待手工操作时也应等待docker compose ps显示数据库健康后再建表。Datacenter 配置cassandra.datacenter必须与目标集群的 seed 配置匹配官方镜像默认数据中心名为datacenter1若自定义集群名请同步修改。Keyspace/表需提前存在RisingWave 的 Cassandra Sink 不会自动建 Keyspace 和表必须先按第三节完成初始化否则 Sink 创建或写入会失败。一致性语义Demo 使用append-only模式若目标表存在与上游主键冲突的数据需要结合业务评估是否改用其他写入策略避免语义不符。至此你已经可以完整复现“RisingWave 流式计算 → 双写 Cassandra / ScyllaDB”的链路先用docker-compose up -d拉起环境再用prepare_cassandra_and_scylladb.sql建好目标库表依次执行三个 SQL 文件建立 Source、物化视图与 Sink最后用 cqlsh 或sink_check.py验证数据落库。该模式同样适用于任何基于 CQL 协议的兼容数据库可平滑迁移到生产环境。赞分享数据库流处理后端数据工程【免费下载链接】risingwaveEvent streaming platform for agentic AI. Continuously ingest, transform, and serve event streams in real time, at scale.项目地址https://gitcode.com/gh_mirrors/ri/risingwave点击查看免费下载相关推荐ScyllaDB CDC Source Connector 完整指南将 ScyllaDB 行级变更实时流式复制到 KafkaScyllaDB CDC Source Connector 完整指南将 ScyllaDB 行级变更实时流式复制到 Kafka 本文围绕 ScyllaDB 官方数据库分布式数据库后端大数据PP-OCRv6-small-det-GGUF技术原理揭秘CrispEmbed优化如何提升检测精度PP OCRv6 small det GGUF技术原理揭秘CrispEmbed优化如何提升检测精度 PP OCRv6 small det GGUF是基于PadScyllaDB 与 Databricks 集成指南基于 Spark Cassandra Connector 的完整实操ScyllaDB 与 Databricks 集成指南基于 Spark Cassandra Connector 的完整实操 本文是一份面向数据工程师与平台开发者数据库分布式数据库后端大数据上一篇终极指南Capybara测试失败自动截图集成CI环境完整方案下一篇告别窗口混乱i3窗口管理器3招恢复默认布局创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

◈

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

◐

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

▲

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。