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

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

发布时间:2026/9/25 3:21:51 来源:云帆数科 栏目:资讯中心
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),仅供参考

相关推荐

如何把 ModelScope 的 AI 模型跑在自己电脑上:10 分钟从安装到出结果
如何把 ModelScope 的 AI 模型跑在自己电脑上:10 分钟从安装到出结果

如何把 ModelScope 的 AI 模型跑在自己电脑上:10 分钟从安装到出结果 【免费下载链接】modelscope ModelScope: bring the notion of Model-as-a-Service to life. 项目地址: https://gitcode.com/GitHub_Trending/mo/modelscope 想跳过云端、在自己的机器上… · 2026/9/25 3:21:51

AlphaFold 3 性能优化实战指南:数据管线、GPU 推理、编译桶与内存配置全解析
AlphaFold 3 性能优化实战指南:数据管线、GPU 推理、编译桶与内存配置全解析

人工智能基础模型深度学习生物信息学科学计算 【免费下载链接】alphafold3 AlphaFold 3 inference pipeline. 项目地址: https://gitcode.com/gh_mirrors/alp/alphafold3 点击查看 免费下载 导读 AlphaFold 3 的推理管线分为数据管线(遗传序列搜索与模… · 2026/9/25 3:21:44

Gixy HTTP Splitting 插件实战:检测 Nginx 配置中的 CRLF 注入与 HTTP 响应拆分漏洞
Gixy HTTP Splitting 插件实战:检测 Nginx 配置中的 CRLF 注入与 HTTP 响应拆分漏洞

静态分析应用安全 【免费下载链接】gixy Nginx configuration static analyzer 项目地址: https://gitcode.com/gh_mirrors/gi/gixy 点击查看 免费下载 导读 HTTP Splitting(HTTP 拆分)是 Nginx 配置中一类常见的输入校验缺陷引发的注入攻击… · 2026/9/25 3:21:38

Apereo CAS OAuth 2.0 授权码流程(Authorization Code)与 PKCE 扩展实战指南
Apereo CAS OAuth 2.0 授权码流程(Authorization Code)与 PKCE 扩展实战指南

后端认证鉴权单点登录 【免费下载链接】cas Apereo CAS - Identity & Single Sign On for all earthlings and beyond. 项目地址: https://gitcode.com/gh_mirrors/ca/cas 点击查看 免费下载 导读 授权码(Authorization Code)是 OAuth … · 2026/9/25 4:26:38

从SQL注入到应急响应:安全工程师面试的闭环答题思路
从SQL注入到应急响应:安全工程师面试的闭环答题思路

每次整理网络安全面试题,我都会提醒候选人:别把希望压在背payload上,真正值钱的答题思路是把“SQL注入怎么发现、怎么防御、出了事怎么应急响应”串成一条闭环。你看标题里“从SQL注入到应急响应”这八个字,其实正是一个安全工程师… · 2026/9/25 4:26:38

ESPnet2 目标说话人提取(TSE)实战:基于 LibriMix 与 TD-SpeakerBeam 的训练、评估与结果解读
ESPnet2 目标说话人提取(TSE)实战:基于 LibriMix 与 TD-SpeakerBeam 的训练、评估与结果解读

人工智能语音音频深度学习NLP 【免费下载链接】espnet End-to-End Speech Processing Toolkit 项目地址: https://gitcode.com/gh_mirrors/es/espnet 点击查看 免费下载 本指南以 ESPnet 仓库中 egs2/librimix/tse1 目标说话人提取(Target Speaker Extr… · 2026/9/25 4:26:32

PS图片出血扩展神器Image Extend:原理、安装与避坑完全指南
PS图片出血扩展神器Image Extend:原理、安装与避坑完全指南

简介:这是一份专为Photoshop设计的图片出血扩展插件Image Extend 1.0.0中文汉化版,面向需要处理印刷品出血位设计的UI设计师、平面设计师及印前工作人员。插件可智能分析图像背景并自动扩展至所需尺寸,支持自定义出血宽度和高度、多图层分别处… · 2026/9/25 4:26:32

AWS HealthImaging 像素数据校验实战:使用 AWS SDK for JavaScript v3 验证 DICOM 解码帧的 CRC32 一致性
AWS HealthImaging 像素数据校验实战:使用 AWS SDK for JavaScript v3 验证 DICOM 解码帧的 CRC32 一致性

示例工程教程后端 【免费下载链接】aws-doc-sdk-examples Welcome to the AWS Code Examples Repository. This repo contains code examples used in the AWS documentation, AWS SDK Developer Guides, and more. For more information, see the Readme.md file below. 项目地… · 2026/9/25 4:26:32

Moto 中的 Bedrock AgentCore 模拟:事件 API 实现与实战指南
Moto 中的 Bedrock AgentCore 模拟:事件 API 实现与实战指南

Mock测试 【免费下载链接】moto A library that allows you to easily mock out tests based on AWS infrastructure. 项目地址: https://gitcode.com/gh_mirrors/mo/moto 点击查看 免费下载 导读 Amazon Bedrock AgentCore 是 AWS 面向智能体(Agent&a… · 2026/9/25 4:26:26

数值优化(Numerical Optimization)学习系列-03-共轭梯度方法(Conjugate Gradient)
数值优化(Numerical Optimization)学习系列-03-共轭梯度方法(Conjugate Gradient)

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 1:00:31

创维E900V22D刷机全攻略:S905L3SB芯片兼容性解析与救砖实战
创维E900V22D刷机全攻略:S905L3SB芯片兼容性解析与救砖实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 1:00:31

MQTT协议原理与Broker服务器搭建实战:从Mosquitto到EMQX
MQTT协议原理与Broker服务器搭建实战:从Mosquitto到EMQX

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 1:00:37

了解更多?预约专属演示

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

企业微信二维码