- 数据库
- 流处理
- 后端
- 数据工程
【免费下载链接】risingwave
Event streaming platform for agentic AI. Continuously ingest, transform, and serve event streams in real time, at scale.
本指南以 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_NAME=cloudinfra 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_NAME=cloudinfra值得注意的是: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, );SimpleStrategy+replication_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 创建 Source:datagen 持续模拟数据
执行 integration_tests/cassandra-and-scylladb-sink/create_source.sql,它创建两张表:
user_behaviors:使用connector = 'datagen'内置连接器,user_id按 1~1000 的序列递增,其余字段随机生成,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_only='true', 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_only='true', cassandra.url = 'scylladb:9042', cassandra.keyspace = 'demo', cassandra.table = 'demo_bhv_table', cassandra.datacenter = 'datacenter1', );5.1 参数逐项说明
| 参数 | 值 | 含义 |
|---|---|---|
connector | cassandra | 指定使用 Cassandra/ScyllaDB 连接器 |
type | append-only | 追加写入模式,不做 upsert 语义 |
force_append_only | true | 强制按 append-only 处理,即使上游可能含更新也忽略其变更语义 |
cassandra.url | cassandra:9042/scylladb:9042 | 目标数据库地址(Docker 网络内服务名:端口) |
cassandra.keyspace | demo | 目标 Keyspace |
cassandra.table | demo_bhv_table | 目标表名 |
cassandra.datacenter | datacenter1 | Cassandra 驱动连接所用的数据中心名,需与集群实际配置一致 |
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→smallint,integer→int,bigint→bigintdecimal→decimal,real→float,double precision→doublevarchar→text,bytea→blobdate→date,time→timetimestamptz→timestamp,interval→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 都监听 9042,docker-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 协议的兼容数据库,可平滑迁移到生产环境。
- 数据库
- 流处理
- 后端
- 数据工程
【免费下载链接】risingwave
Event streaming platform for agentic AI. Continuously ingest, transform, and serve event streams in real time, at scale.
相关推荐
ScyllaDB CDC Source Connector 完整指南:将 ScyllaDB 行级变更实时流式复制到 Kafka
ScyllaDB CDC Source Connector 完整指南:将 ScyllaDB 行级变更实时流式复制到 Kafka 本文围绕 ScyllaDB 官方
数据库分布式数据库后端大数据PP-OCRv6-small-det-GGUF技术原理揭秘:CrispEmbed优化如何提升检测精度
PP OCRv6 small det GGUF技术原理揭秘:CrispEmbed优化如何提升检测精度 PP OCRv6 small det GGUF是基于Pad
ScyllaDB 与 Databricks 集成指南:基于 Spark Cassandra Connector 的完整实操
ScyllaDB 与 Databricks 集成指南:基于 Spark Cassandra Connector 的完整实操 本文是一份面向数据工程师与平台开发者
数据库分布式数据库后端大数据
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考