
DataHub 以 Cassandra 作为元数据存储后端本地快速启动与源码级原理解析【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub本篇文章围绕 docker/cassandra/README.md 展开讲解如何让 DataHub 的元数据服务 GMS 使用 Apache Cassandra 作为替代存储后端从零完成本地 quickstart 实例的构建、启动与示例数据摄入并结合仓库内 Compose 覆盖文件、CQL 初始化脚本与metadata-io模块中的 Cassandra DAO 源码深入说明 Cassandra 侧的表结构设计、连接配置与读写调用链。读完本文你将掌握一套可复现的GMS Cassandra本地部署方案并能从实现层理解 DataHub 如何以metadata_aspect_v2表承载全部元数据。一、为什么 DataHub 需要 Cassandra 存储后端DataHub 的元数据服务 GMSGeneralized Metadata Service在默认部署中依赖关系型数据库如 MySQL/PostgreSQL通过 Ebean 持久化作为主存储。但 DataHub 的核心元数据模型本质上是实体Entity 方面Aspect的键值结构这与 Cassandra 的宽列模型天然契合。因此DataHub 在metadata-io模块中抽象出AspectDao接口并提供了基于 DataStax Java Driver 的 Cassandra 实现让 GMS 可以切换到 Cassandra 作为主存储。从 docker/datahub-gms/env/docker.cassandra.env 可以看到切换后端的关键开关只有一个环境变量ENTITY_SERVICE_IMPLcassandra当 GMS 以该配置启动时实体方面aspect的读写将由CassandraAspectDao承担其余组件Kafka、Elasticsearch、Neo4j 图服务保持不变。这正是 docker/cassandra/README.md 所述 DataHub GMS can use Cassandra as an alternate storage backend 的工程基础。二、环境准备与源码构建2.1 用 Gradle 构建项目从源码启动的前提是先把项目构建出来。在仓库根目录执行./gradlew build2.2 解决 Java 版本不兼容问题构建时如果遇到如下报错Required org.gradle.jvm.version 8 and found incompatible value 11.说明当前默认 JVM 版本与 Gradle 所需版本不匹配DataHub 的 Gradle 构建链依赖 Java 8 工具链。此时需要先把JAVA_HOME指向 Java 8 再重试export JAVA_HOME/usr/lib/jvm/java-8-openjdk-amd64 $JAVA_HOME/bin/java -version用java -version确认切换生效后重新执行./gradlew build即可继续。该步骤来自原文档属于本地源码构建的常见前置坑点。三、以 Cassandra 为后端启动 DataHub3.1 一键脚本启动构建完成后原文档给出的启动方式是在仓库根目录执行./docker/dev-with-cassandra.sh该脚本的作用是从源码构建并启动所有 DataHub 组件且 GMS 使用 Cassandra 后端。需要说明的是当前仓库快照中并未收录该脚本文件其等效的声明式启动方式可以直接使用docker目录下的 Compose 覆盖文件完成见下文 3.2效果一致。3.2 使用 Compose 覆盖文件启动当前仓库内的等效方案当前仓库在 docker/cassandra/docker-compose.cassandra.yml同一份副本也维护在 docker/profiles/cassandra/docker-compose.cassandra.yml中提供了完整的 Cassandra 后端覆盖定义。将其与主 Compose 文件叠加即可启动services: cassandra: hostname: cassandra image: cassandra:3.11 ports: - 9042:9042 healthcheck: test: [CMD, cqlsh, -u cassandra, -p cassandra, -e describe keyspaces] interval: 15s timeout: 10s retries: 10 volumes: - cassandradata:/var/lib/cassandra cassandra-load-keyspace: image: cassandra:3.11 depends_on: cassandra: condition: service_healthy volumes: - ./cassandra/init.cql:/init.cql command: /bin/bash -c cqlsh cassandra -f /init.cql datahub-gms: env_file: datahub-gms/env/docker.cassandra.env depends_on: - cassandra volumes: cassandradata:这个覆盖文件包含三个关键设计Cassandra 节点使用官方cassandra:3.11镜像宿主机映射9042CQL 原生协议端口通过cqlsh -u cassandra -p cassandra -e describe keyspaces做健康检查15 秒间隔、10 次重试只有 keyspace 可查询才算就绪。一次性 keyspace 初始化容器cassandra-load-keyspace在 Cassandra 健康后挂载并执行 docker/cassandra/init.cql完成 keyspace 与元数据表结构创建详见第四节。GMS 环境注入datahub-gms通过env_file: datahub-gms/env/docker.cassandra.env注入 Cassandra 连接参数并在depends_on中声明对 Cassandra 的依赖。数据卷cassandradata将 Cassandra 数据持久化到/var/lib/cassandra重启服务不会丢失已摄入的元数据。四、Cassandra 侧的元数据模型keyspace 与 metadata_aspect_v24.1 初始化脚本完整解析docker/cassandra/init.cql 是 Cassandra 后端的建库建表脚本其全部内容如下create keyspace if not exists datahub WITH REPLICATION { class : SimpleStrategy, replication_factor : 1 }; create table if not exists datahub.metadata_aspect_v2 ( urn varchar, aspect varchar, version bigint, metadata text, systemmetadata text, createdon timestamp, createdby varchar, createdfor varchar, entity varchar, primary key ((urn), aspect, version)) with clustering order by (aspect asc, version asc); insert into datahub.metadata_aspect_v2 (urn, aspect, version, metadata, systemmetadata, createdon, createdby, entity) values( urn:li:corpuser:datahub, corpUserInfo, 0, {displayName:Data Hub,active:true,fullName:Data Hub,email:datahublinkedin.com}, {}, toTimestamp(now()), urn:li:corpuser:__datahub_system, corpuser ) if not exists; insert into datahub.metadata_aspect_v2 (urn, aspect, version, metadata, systemmetadata, createdon, createdby, entity) values( urn:li:corpuser:datahub, corpUserEditableInfo, 0, {skills:[],teams:[],pictureLink:https://raw.githubusercontent.com/linkedin/datahub/master/datahub-web-react/src/images/default_avatar.png}, {}, toTimestamp(now()), urn:li:corpuser:__datahub_system, corpuser ) if not exists;脚本依次完成三件事创建 keyspacedatahubkeyspace 使用SimpleStrategy、副本因子 1适用于本地单节点 quickstart生产多机房环境应替换为NetworkTopologyStrategy并相应调大副本因子。创建核心表metadata_aspect_v2这是 DataHub 在 Cassandra 中的唯一元数据事实表表名常量定义见CassandraAspect.TABLE_NAME其复合主键设计为分区键urn实体 URN聚簇键aspect方面名version版本号聚簇顺序(aspect asc, version asc)这意味着同一实体的所有方面按名称排序、同一方面按版本升序排列天然支持按 URN 读取实体的全部方面以及读取某一方面的最新版本两类高频查询。其余列中metadata存放方面内容的 JSON 文本systemmetadata存放系统级元数据如运行 ID、系统运行标识等createdon/createdby/createdfor记录审计信息entity记录实体类型如corpuser。写入内置用户种子数据预置urn:li:corpuser:datahub用户的corpUserInfo与corpUserEditableInfo两个方面保证系统启动后即存在默认登录用户datahub密码在初始化流程中对应为datahub。两条语句均带if not exists可安全重复执行。4.2 表结构与 Java 实体一一对应Cassandra 侧的表结构与 GMS 的 Java 实体在 CassandraAspect.java 中严格对齐字段urn、aspect、version、metadata、systemMetadata、createdOn、createdBy、createdFor与表列一一对应常量TABLE_NAME metadata_aspect_v2与 init.cql 中的表名一致静态方法rowToEntityAspect(Row row)负责把 Cassandra 查询结果Row转换为内存共享表示EntityAspect是存储层与上层模型之间的转换枢纽。从实现上看CassandraAspect的列名常量URN_COLUMN、ASPECT_COLUMN、VERSION_COLUMN等与 init.cql 中列名完全一致任何一端改动列名都会导致另一端读写失败因此两者必须同步维护。五、GMS 的 Cassandra 连接配置详解GMS 侧的全部 Cassandra 相关配置集中在 docker/datahub-gms/env/docker.cassandra.envDATAHUB_UPGRADE_HISTORY_KAFKA_CONSUMER_GROUP_IDgeneric-duhe-consumer-job-client-gms KAFKA_BOOTSTRAP_SERVERbroker:29092 KAFKA_SCHEMAREGISTRY_URLhttp://schema-registry:8081 ELASTICSEARCH_HOSTelasticsearch ELASTICSEARCH_PORT9200 ES_BULK_REFRESH_POLICYNONE ELASTICSEARCH_INDEX_BUILDER_SETTINGS_REINDEXtrue ELASTICSEARCH_INDEX_BUILDER_MAPPINGS_REINDEXtrue NEO4J_HOSThttp://neo4j:7474 NEO4J_URIbolt://neo4j NEO4J_USERNAMEneo4j NEO4J_PASSWORDdatahub JAVA_OPTS-Xms1g -Xmx1g GRAPH_SERVICE_IMPLneo4j ENTITY_REGISTRY_CONFIG_PATH/datahub/datahub-gms/resources/entity-registry.yml MAE_CONSUMER_ENABLEDtrue MCE_CONSUMER_ENABLEDtrue CASSANDRA_DATASOURCE_USERNAMEcassandra CASSANDRA_DATASOURCE_PASSWORDcassandra CASSANDRA_HOSTScassandra CASSANDRA_PORT9042 CASSANDRA_DATASOURCE_HOSTcassandra:9042 ENTITY_SERVICE_IMPLcassandra其中与 Cassandra 后端直接相关的核心参数如下环境变量取值说明ENTITY_SERVICE_IMPLcassandra切换存储后端的总开关决定 GMS 使用CassandraAspectDao而非 Ebean 实现CASSANDRA_HOSTScassandraCassandra 节点主机名Compose 网络内服务名CASSANDRA_PORT9042CQL 原生协议端口CASSANDRA_DATASOURCE_HOSTcassandra:9042连接串形式的主机与端口CASSANDRA_DATASOURCE_USERNAMEcassandra认证用户名与 Compose 健康检查、init.cql 执行所用的凭据一致CASSANDRA_DATASOURCE_PASSWORDcassandra认证密码其余变量用于维持 DataHub 其余组件的正常运行KAFKA_BOOTSTRAP_SERVER指向 Kafka brokerKAFKA_SCHEMAREGISTRY_URL指向 Schema RegistryELASTICSEARCH_*系列指向搜索索引NEO4J_*与GRAPH_SERVICE_IMPLneo4j指示图服务仍由 Neo4j 承担MAE_CONSUMER_ENABLED/MCE_CONSUMER_ENABLED打开元数据事件消费链路。也就是说Cassandra 替换的仅仅是实体方面aspect的主存储搜索、图与事件链路保持不变。六、启动后摄入示例数据所有服务就绪后在另一个终端进入源码项目自带的 Python 虚拟环境初始化 CLI 并摄入示例数据source metadata-ingestion/venv/bin/activate datahub init --username datahub --password datahub datahub datapack load showcase-ecommercedatahub init用datahub/datahub建立 CLI 与 GMS 的连接对应第四节种子数据中预置的urn:li:corpuser:datahub用户datahub datapack load showcase-ecommerce加载showcase-ecommerce数据包向 DataHub 摄入一整套电商场景的演示元数据数据集、图表、仪表盘、血缘等。摄入完成后即可通过 DataHub 前端检索这批元数据并验证它们在 Cassandra 的datahub.metadata_aspect_v2表中以 (urn, aspect, version) 形式落盘。七、源码级原理Cassandra 读写链路7.1 DAO 层实现Cassandra 后端的读写入口是 CassandraAspectDao.java它同时实现AspectDao与AspectMigrationsDao两个接口因此既负责正常读写也承担索引重建、数据迁移等运维操作。其关键点包括构造函数通过PrimaryStorageResolver.resolveCassandraPrimary()解析出CqlSessionDataStax Java Driver 的会话对象并在构造时注入系统方面校验器SystemAspectValidator与方面大小校验配置validateConnection()启动时校验若metadata_aspect_v2表不存在则记录日志 GMS cant find entity aspects table in Cassandra storage layer. 并将canWrite置为false——这正是 init.cql 必须在 GMS 启动前执行完毕的原因getLatestAspect()按(urn, aspect, ASPECT_LATEST_VERSION)读取最新版本方面并在 forUpdate 场景下先执行预补丁校验再包装为SystemAspect供上层使用类中还引用了PartitionedStream、OffsetPager等分页机制用于支撑大规模列表查询与索引回填时的流式扫描。7.2 启动自检表存在性检查实现在 AspectStorageValidationUtil.java 的checkTableExists(CqlSession)方法中通过查询系统表system_schema.tables对应 SQLSELECT table_name ...的 CQL 等价物判断目标表是否已建立。这解释了快速启动流程中先起 Cassandra、执行 init.cql、再起 GMS的严格顺序约束。7.3 生命周期管理Cassandra 后端还包含对应的保留策略服务 CassandraRetentionService.java用于按时限清理过期方面版本对应 DataHub 的方面保留机制使 Cassandra 后端与默认后端在数据治理能力上对齐。八、小结以 Cassandra 作为 DataHub GMS 的替代存储后端是一条完整且可复现的本地部署路径./gradlew build完成构建dev-with-cassandra.sh或 docker/cassandra/docker-compose.cassandra.yml 负责拉起全部组件docker/cassandra/init.cql 建立datahubkeyspace 与metadata_aspect_v2表并写入种子用户ENTITY_SERVICE_IMPLcassandra与 docker/datahub-gms/env/docker.cassandra.env 中的连接参数完成 GMS 侧切换最后通过datahub initdatahub datapack load showcase-ecommerce验证端到端可用。从源码看CassandraAspectDao与CassandraAspect共同构成了这套后端的读写核心其表结构与 init.cql 严格一一对应理解这一对应关系是后续自建 Cassandra 集群、调整副本策略或排查启动失败如表未初始化时最有效的抓手。【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考