数据库流处理后端数据工程【免费下载链接】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 的 Meta 服务使用关系型数据库PostgreSQL / MySQL / SQLite持久化系统目录Catalog与集群元数据这些持久化层的 Rust 实体定义集中在src/meta/model下并由src/meta/model/migration这个独立的 SeaORM 迁移工程负责版本化演进。本文以 src/meta/model/src/README.md 为核心骨架完整讲解定义版本间变更 → 生成迁移文件 → 应用到数据库 → 反向生成模型文件的标准化工作流并深入剖析仓库中真实迁移脚本、派生宏与后端差异处理帮助你掌握在 RisingWave 元数据模型上安全增删表、变更列、维护枚举与数组类型的全套实战方法。一、认识 RisingWave 的 Meta 数据模型与迁移工程在 RisingWave 中Meta 节点负责维护集群的 Catalog库、表、物化视图、源、Sink、索引、函数、用户与权限等以及 Hummock 存储层、流作业调度等运行态元数据。这些数据需要落库持久化其 ORM 层采用 SeaORM 框架实体Entity与迁移Migration分别组织在两个 crate 中目录角色src/meta/model/srcrisingwave_meta_modelcrate存放全部 SeaORM 实体Entity与辅助宏src/meta/model/migrationrisingwave_meta_model_migrationcrateSeaORM Migrator CLI 与全部迁移文件src/meta/model/tests实体与迁移的一致性测试从源码结构看model/src目前包含 47 个实体模块以 lib.rs 中的pub mod声明为准覆盖database、schema、table、source、sink、index、function、user、worker、streaming_job、fragment、hummock_version_delta、object_dependency等核心领域对象而migration/src中从m20230908_072257_init到m20260805_000000_object_belong_to_oid已累计注册 74 个迁移文件见 lib.rs。src/meta/model/src/README.md给出的核心思路是用迁移脚本驱动数据库 Schema 演进再从数据库反向生成模型文件避免手工维护实体代码同时把 PostgreSQL 特有的数组、枚举等类型单独手工定义在模型文件中。二、版本变更的标准工作流三步骤总览原文档将一次元数据模型变更归纳为如下流程生成新迁移文件并应用到数据库使用 Migrator CLI 的generate与up子命令得到一个可执行的迁移脚本并真正改变数据库结构从数据库反向生成模型文件使用sea-orm-cli generate entity读取最新表结构自动生成实体代码拷贝到模型 crate 中手工补充 PG 专属类型与辅助函数数组、枚举等类型需要手工定义在模型文件中必要时再补充辅助函数。下面逐节展开每个步骤的命令、代码写法与仓库内的真实范例。三、第一步生成新迁移文件在 src/meta/model/src/README.md 中生成迁移文件以本地 PostgreSQL 为例export DATABASE_URLpostgres://postgres:localhost:5432/postgres cargo run -- generate MIGRATION_NAME cargo run -- up需要说明的是MIGRATION_NAME是一个描述性标识符Migrator 会根据它生成一个形如m20240101_000000_migration_name.rs的独立文件必须在src/meta/model/migration目录下运行而不是项目根目录migration/README.md 对此有明确提示DATABASE_URL是必填的连接串但generate子命令需要数据库端点但实际上并不会使用它——也就是说生成阶段即使数据库不可用也可以完成要真正应用迁移时才需要可用的数据库。main.rs 展示了 CLI 的入口实现sea_orm_migration::prelude::cli::run_cli(risingwave_meta_model_migration::Migrator)。这里传入的Migrator类型定义在 lib.rs 中它的migrations()方法按时间顺序返回一个VecBoxdyn MigrationTrait即全部已注册迁移的调度表——新增迁移文件后记得在这里追加一行Box::new(module::Migration)否则 Migrator 不会执行它。迁移文件命名与编号从仓库中已有的迁移文件可以看出命名规范为mYYYYMMDD_HHMMSS_snake_case_name.rs例如m20230908_072257_init.rs—— 初始 Schemam20240617_070131_index_column_properties.rs—— 给index表加列README 中推荐的参考样例m20250106_072104_fragment_relation.rs—— 新增 fragment 关联表m20250810_000000_add_user_admin_field.rs—— 给用户表加管理员字段m20260805_000000_object_belong_to_oid.rs—— 最近期的幂等迁移文件名中的时间戳决定了迁移的执行顺序SeaORM 会为每个已执行的迁移在数据库中登记版本记录保证每个迁移文件只被应用一次。四、第二步编写迁移脚本up / downgenerate生成的文件会包含一份模板迁移脚本需要你替换为自己的实现。原文档给出了MigrationTrait的标准骨架#[async_trait::async_trait] impl MigrationTrait for Migration { async fn up(self, manager: SchemaManager) - Result(), DbErr { // Replace the sample below with your own migration scripts todo!(); } async fn down(self, manager: SchemaManager) - Result(), DbErr { // Replace the sample below with your own migration scripts todo!(); } }up定义升级操作down定义对应的回滚操作。原文档指出你可以在迁移文件中定义表、索引、外键也可以执行任意 DML 操作来保证数据正确性例如数据回填、格式转换。4.1 一个最小可用的增列迁移原文档推荐的参考样例 m20240617_070131_index_column_properties.rs 非常典型#[async_trait::async_trait] impl MigrationTrait for Migration { async fn up(self, manager: SchemaManager) - Result(), DbErr { manager .alter_table( Table::alter() .table(Index::Table) .add_column(ColumnDef::new(Index::IndexColumnProperties).rw_binary(manager)) .to_owned(), ) .await } async fn down(self, manager: SchemaManager) - Result(), DbErr { manager .alter_table( Table::alter() .table(Index::Table) .drop_column(Index::IndexColumnProperties) .to_owned(), ) .await } } #[derive(DeriveIden)] enum Index { Table, IndexColumnProperties, }这里有两个值得学习的点每个被引用的表名、列名都要用#[derive(DeriveIden)]枚举声明SeaORM 会将其转换为合法标识符避免手写字符串带来的大小写与保留字问题大字段列使用rw_binary(manager)而不是binary()这是仓库自带的跨后端扩展见下文第六节。4.2 建表、索引与外键以 init 迁移为例初始化迁移 m20230908_072257_init.rs 是仓库中规模最大的迁移它完整展示了 SeaORM 建表 API 的用法也定义了 RisingWave 元数据模型的基本盘Cluster 表以 UUID 作为集群 IDWorker / WorkerProperty 表Worker 表记录节点host、port、status、worker_typeWorkerProperty 通过外键FK_worker_property_worker_id关联 Worker并带ON DELETE CASCADEUser / Object / ObjectDependency / UserPrivilege 表构成完整的用户-对象-权限模型Object表自引用schema_id、database_id都指向Object.oidDatabase / Schema / StreamingJob / Fragment / Actor / ActorDispatcher 表覆盖 Catalog 与流图执行元数据Fragment.StreamNode、Actor.Splits、Actor.ExprContext等大字段均用rw_binary存储 protobuf 序列化数据Table / Source / Sink / Index / View / Function 表完整刻画表、源、Sink、索引、视图、函数对象其中DefinitionSQL 定义文本用rw_long_textSystemParameter / CatalogVersion 表系统参数与目录版本号。索引与外键的创建方式manager .create_index( MigrationIndex::create() .table(Worker::Table) .name(idx_worker_host_port) .unique() .col(Worker::Host) .col(Worker::Port) .to_owned(), ) .await?; manager .create_table( MigrationTable::create() .table(WorkerProperty::Table) .col(ColumnDef::new(WorkerProperty::WorkerId).integer().primary_key()) // ... .foreign_key( mut ForeignKey::create() .name(FK_worker_property_worker_id) .from(WorkerProperty::Table, WorkerProperty::WorkerId) .to(Worker::Table, Worker::WorkerId) .on_delete(ForeignKeyAction::Cascade) .to_owned(), ) .to_owned(), ) .await?;init 迁移还演示了初始化数据的写法通过Query::insert()创建集群 ID随机 UUID、内置用户root/postgres/rwadmin、内置数据库dev以及public/pg_catalog/information_schema/rw_catalog四个内置 Schema并在 MySQL 与 PostgreSQL 下用ALTER TABLE ... AUTO_INCREMENT/SELECT setval(...)重置自增序列起点保证后续对象的 OID 从固定值开始。4.3 不要修改已发布的历史迁移migration/README.md 的警告框强调了一个重要约束每个迁移文件只能被应用一次并被记录在系统表中。对于新的 Schema 变更必须生成新的迁移文件。除非你确信对迁移文件的修改还没有包含在任何已发布版本中否则不要修改已经发布的迁移文件。这正是版本演进的意义已发布的迁移已进入线上数据库的版本记录改动它会导致不同环境间 Schema 不一致。遇到需要调整的 Schema一律通过新的迁移文件来追加变更。五、第三步应用迁移并查看状态migration/README.md 完整列出了 Migrator CLI 支持的子命令整理如下子命令作用cargo run -- generate MIGRATION_NAME生成新的迁移文件需要DATABASE_URL但不实际连接使用cargo run/cargo run -- up应用所有待执行的迁移cargo run -- up -n 10只应用前 10 个待执行的迁移cargo run -- down回滚最近一次应用的迁移cargo run -- down -n 10回滚最近 10 次应用的迁移cargo run -- fresh删除数据库全部表然后重新应用所有迁移cargo run -- refresh回滚所有已应用的迁移再重新应用所有迁移cargo run -- reset回滚所有已应用的迁移不重新应用cargo run -- status查看所有迁移的状态在开发调试阶段还可以用 SQLite 内存数据库快速验证迁移逻辑而无需启动真实的 PostgreSQLDATABASE_URLsqlite::memory: cargo run -- generate MIGRATION_NAME需要特别留意的是 MySQL 后端下的约束Migrator 只有在up返回后才记录迁移版本而MySQL 的 DDL 会隐式提交事务。这意味着执行中途崩溃会导致迁移半途执行但版本未记录。仓库中较新的迁移如 m20260805_000000_object_belong_to_oid.rs会刻意把每个 DDL 步骤包在has_column/has_index/has_foreign_key等存在性检查中使新的 Meta 主节点可以安全地重试一个只执行了一部分的迁移——这是编写高可用环境迁移脚本的重要模式。六、第四步从数据库生成模型文件迁移应用完毕、表结构确定之后就可以用sea-orm-cli反向生成实体模型了。原文档给出的命令如下cargo run -- up sea-orm-cli generate entity -u postgres://postgres:localhost:5432/postgres -s public -o {target_dir} cp {target_dir}/xxx.rs src/meta/src/model/对这条命令的解读与补充-u指定数据库连接串-s指定 schemaPostgreSQL 场景下通常为public-o指定输出目录生成出的实体代码无需手工编写拷贝到模型 crate 后还要完成三件收尾工作在 src/meta/model/src/lib.rs 中pub mod声明该模块如果它是新的实体表把它加入for_all_meta_model_entities!宏的实体清单按需在 prelude.rs 中重新导出Entity方便业务代码统一use。注原文档中拷贝目标写为src/meta/src/model/在本仓库的实际布局中模型文件的目标目录应为src/meta/model/src/即 src/meta/model/src请以实际目录为准。实体清单的一致性保障lib.rs中的for_all_meta_model_entities!宏用macro_rules维护了一份全量实体清单lib.rs而 tests/meta_model_entities.rs 中的测试会用syn解析src/下所有.rs文件提取每个#[sea_orm(table_name ...)]声明的表名与for_all_meta_model_entities!枚举出的模块逐一比对断言两者完全一致否则报出Missing in for_all_meta_model_entities / Unexpected in for_all_meta_model_entities。这意味着新增实体表后必须同步更新宏清单否则 CI 中的该测试会失败。七、手工定义 PG 专属类型数组与枚举原文档强调数组Array与枚举Enum类型基本只有 PostgreSQL 原生支持因此需要在模型文件中手工定义。仓库中的做法不是直接暴露原生类型而是借助 SeaORM 的FromJsonQueryResult/DeriveActiveEnum等派生宏包装成可落库、可序列化的 Rust 类型。7.1 数组字段包装为 JSON 类型原文档给出的I32Array示例// We define integer array typed fields as json and derive it using the follow one. #[derive(Clone, Debug, PartialEq, FromJsonQueryResult, Eq, Serialize, Deserialize, Default)] pub struct I32Array(pub Veci32);这样整数数组在数据库中以 JSON 形式存储SeaORM 的FromJsonQueryResult派生负责 JSON 与 Rust 结构之间的互相转换。该类型在 lib.rs 中通过仓库自带的derive_from_json_struct!宏实例化derive_from_json_struct!(TableIdArray, VecTableId); derive_from_json_struct!(EpochArray, VecEpoch); derive_from_json_struct!(I32Array, Veci32); derive_from_json_struct!(Property, BTreeMapString, String);其中Property被大量用于存储表 / Source / Sink 的 with 属性键值对。7.2 枚举字段映射为字符串原文档给出的枚举示例// We define enum typed fields as string and derive it using the follow one. #[derive(Clone, Debug, PartialEq, Eq, EnumIter, DeriveActiveEnum)] #[sea_orm(rs_type String, db_type String(None))] pub enum WorkerStatus { #[sea_orm(string_value STARTING)] Starting, #[sea_orm(string_value RUNNING)] Running, }关键点在于rs_type String表示 Rust 侧类型为字符串db_type String(None)表示数据库中存储为字符串列#[sea_orm(string_value ...)]显式指定每个枚举变体对应的数据库字符串值。仓库中实际使用这一模式的枚举包括 lib.rs 中的JobStatusINITIAL/CREATING/CREATED、CreateTypeBACKGROUND/FOREGROUND与DispatcherTypeHASH/BROADCAST/SIMPLE/NO_SHUFFLE它们都额外实现了与 protobuf 枚举PbStreamJobStatus、PbCreateType、PbDispatcherType的双向转换保证数据库字符串值与 gRPC 协议枚举一一对应。例如CreateType的转换impl FromCreateType for PbCreateType { fn from(create_type: CreateType) - Self { match create_type { CreateType::Background Self::Background, CreateType::Foreground Self::Foreground, } } }StreamingParallelismAdaptive/Fixed(usize)/Custom则作为一个 JSON 枚举直接存储在streaming_job表的parallelismJSON 列中展示了枚举复杂载荷的另一种表达方式。7.3 复杂载荷protobuf 二进制包装宏除了数组与枚举RisingWave 的元数据表中还大量存储 protobuf 序列化后的结构流节点、列目录、SST 信息等。为此 lib.rs 提供了三个底层宏宏用途实例derive_from_json_struct!包装 JSON 存储字段TableIdArray、EpochArray、I32Array、Propertyderive_from_blob!包装单个 protobuf 二进制DeriveValueTypeStreamNode、DataType、ColumnCatalog、TableVersion、SecretRef、WorkerResource等derive_array_from_blob!包装 protobuf 对象数组DataTypeArray、FieldArray、ColumnCatalogArray、HummockVersionDeltaArray、SstableInfoArray等derive_btreemap_from_blob!包装 protobuf BTreeMapSecretRefBTreeMapString, PbSecretRef例如derive_from_blob!展开出的结构以Vecu8存列并提供to_protobuf()/from_protobuf()在 protobuf 消息与存储字节之间互转macro_rules! derive_from_blob { ($struct_name:ident, $field_type:ty) { #[derive(Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize, DeriveValueType)] pub struct $struct_name(#[sea_orm] Vecu8); impl $struct_name { pub fn to_protobuf(self) - $field_type { prost::Message::decode(self.0.as_slice()).unwrap() } fn from_protobuf(val: $field_type) - Self { Self(prost::Message::encode_to_vec(val)) } } // ... }; }理解这几个宏就能看懂table.rs、fragment.rs等实体文件中那些长得不像普通列的字段类型也方便为新的复杂字段选择正确的包装方式。八、跨后端差异MySQL 的长度限制与ColumnDefExt原文档及其引用的 migration/README.md特别强调了一个 MySQL 后端陷阱MySQL 的VARCHAR、TEXT、BLOB、BINARY等类型相对 PostgreSQL / SQLite 有更严格的最大长度限制最大 65,535 字节。当需要存储更大的数据如 SQL 定义、UDF body、protobuf 编码的内部数据等时避免使用ColumnDef::text或ColumnDef::blob这类内置构造器改用 ./src/utils.rs 中定义的扩展方法。对应的实现位于 migration/src/utils.rs它通过easy_ext给ColumnDef增加了两个方法/// Set column type as longblob for MySQL, bytea for Postgres, and blob for Sqlite. pub fn rw_binary(mut self, manager: SchemaManager) - mut Self { match manager.get_database_backend() { DatabaseBackend::MySql self.custom(extension::mysql::MySqlType::LongBlob), DatabaseBackend::Postgres | DatabaseBackend::Sqlite self.blob(), } } /// Set column type as longtext for MySQL, and text for Postgres and Sqlite. pub fn rw_long_text(mut self, manager: SchemaManager) - mut Self { match manager.get_database_backend() { DatabaseBackend::MySql self.custom(Alias::new(longtext)), DatabaseBackend::Postgres | DatabaseBackend::Sqlite self.text(), } }规则总结用途推荐方法MySQLPostgreSQL / SQLite大二进制protobuf 序列化等rw_binary(manager)longblobbytea/blob大文本SQL 定义、UDF bodyrw_long_text(manager)longtexttext在 init 迁移中还能看到另一处 MySQL 特殊处理MySQL 默认以utf8_general_ci编码字符串大小写不敏感而 RisingWave 需要大小写敏感的比较因此在建表前会执行ALTER DATABASE CHARACTER SET utf8mb4 COLLATE utf8mb4_bin调整数据库排序规则m20230908_072257_init.rs。九、实战一个完整的幂等迁移 数据回填范例如果要为线上集群新增一列并回填历史数据仓库中最新、也最值得借鉴的样例是 m20260805_000000_object_belong_to_oid.rs。它给object表添加belong_to_oid列记录对象归属的父对象并展示了三个高级模式1. 逐步骤幂等MySQL DDL 隐式提交场景if !manager.has_column(object, belong_to_oid).await? { // ... add column } if !manager.has_index(object, INDEX_NAME).await? { // ... create index } if matches!(backend, DatabaseBackend::MySql | DatabaseBackend::Postgres) !has_foreign_key(manager).await? { // ... add foreign key }每个 DDL 步骤前都做存在性检查保证迁移在部分执行后可以被安全重试。2. 按后端分写的 SQL 回填由于 SQLite 无法单独为既有表追加外键它在ALTER TABLE时直接以内联REFERENCES的方式定义外键而后端差异更大的回填 UPDATE 语句则分别为 MySQL / PostgreSQL / SQLite 各写一份如belongs_to_job_id回填、__iceberg_sink_前缀隐式对象的归属推导等保持三端行为一致。3. 内嵌单元测试该文件末尾自带两个#[tokio::test]test_sqlite_backfill_and_cascade在sqlite::memory:上构造迷你表结构执行迁移后断言belong_to_oid回填结果与级联删除行为test_sqlite_partial_run_retry则模拟Meta 在第一个 DDL 语句后崩溃但未记录版本与全部语句执行完但未记录版本两种场景验证迁移可重复执行且结果幂等。这为如何给迁移写测试提供了直接范本。十、模型文件中的辅助函数与转换逻辑原文档最后一条是如有必要在模型文件中定义其他辅助函数。仓库中的典型用法包括枚举与 protobuf 的双向转换如前文CreateType、JobStatus、DispatcherType与Pb*类型的From实现lib.rs、lib.rs供控制器controller层在数据库枚举与 gRPC 协议之间互转包装类型的便捷方法如I32Array的into_u32_array()把Veci32转成Vecu32见 lib.rs以及derive_from_blob!生成的to_protobuf()/inner_ref()等预置类型别名TransactionId i32、Epoch i64、CompactionTaskId i64等语义化别名lib.rs让实体字段的意图更清晰。这些辅助函数通常服务于src/meta下的 controller 与 service 层是模型文件 — 业务逻辑之间的衔接层。十一、小结一次完整变更的检查清单综合原文档与仓库源码一次 RisingWave 元数据模型变更应完成以下动作在 src/meta/model/migration 目录下执行DATABASE_URL... cargo run -- generate MIGRATION_NAME生成迁移文件在迁移文件中实现up/down定义或修改表、索引、外键必要时执行数据 DML大字段用rw_binary/rw_long_text在 migration/src/lib.rs 的migrations()中注册新迁移执行cargo run -- up开发期可用fresh/refresh/status调试确认 Schema 与数据正确用sea-orm-cli generate entity从数据库生成实体拷贝到 src/meta/model/src在 lib.rs 中声明模块、加入for_all_meta_model_entities!清单必要时在 prelude.rs 中导出 Entity数组、枚举等 PG 专属类型手工定义为FromJsonQueryResult/DeriveActiveEnum包装类型protobuf 复杂载荷用derive_from_blob!系列宏为复杂迁移编写幂等保护与单元测试参考 m20260805_000000_object_belong_to_oid.rs 的测试写法并保证 meta_model_entities.rs 的一致性测试通过记住红线已发布的迁移文件不可修改新变更一律追加新迁移。按照这套流程你可以在不破坏线上数据的前提下安全、可回滚、可测试地推动 RisingWave 元数据模型的持续演进。赞分享数据库流处理后端数据工程【免费下载链接】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点击查看免费下载相关推荐Loco 模型实战指南基于 SeaORM 的 ActiveRecord 建模、迁移与测试Loco 模型实战指南基于 SeaORM 的 ActiveRecord 建模、迁移与测试 本文围绕 Loco 框架Rust中 Models 一节的完整后端Karakeep 数据库迁移实战基于 Drizzle ORM 的 Schema 演进、迁移生成与 Drizzle Studio 操作指南Karakeep 数据库迁移实战基于 Drizzle ORM 的 Schema 演进、迁移生成与 Drizzle Studio 操作指南 本篇技术指南聚焦当前后端前端移动开发AI 应用知识管理全文检索MCP 服务用 CopilotKit LangGraph Tavily 构建具备 Human-in-the-Loop 能力的 Agent 研究画布应用open-research-ANA用 CopilotKit LangGraph Tavily 构建具备 Human in the Loop 能力的 Agent 研究画布应用open r人工智能AI AgentAgent 框架前端后端上一篇零基础玩转LangChain文件管理Agent从工具链到自动化实战下一篇Cube Studio项目中的Containerd安装与配置指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考