cassandra-data-migrator (also known as CDM)
在源和目标 Cassandra 集群之间迁移和验证表。
[!IMPORTANT] 请注意以下列出的先决条件,以避免 CDM、Spark、Scala、Hadoop 版本不匹配问题
作为容器安装
- 从 DockerHub 获取包含所有依赖项的最新镜像
- 所有迁移工具(
cassandra-data-migrator+dsbulk+cqlsh)将在容器的/assets/文件夹中可用
- 所有迁移工具(
作为 JAR 文件安装
- 使用以下方法之一下载最新的 jar 文件
,
curl -L0 https://github.com/datastax/cassandra-data-migrator/releases/download/x.y.z/cassandra-data-migrator-x.y.z.jar --output cassandra-data-migrator-x.y.z.jar(OR)- 从 packages area here 下载
CDM 6.x+ 先决条件
- Java17(最低要求),因为 Spark 4.x 二进制文件是使用它编译的。
- 如有需要,可以使用任何 Java 17+ LTS 版本(例如 Java 21 或 Java 25)
- Spark 版本 4.1.2
可以通过运行以下命令安装 Spark: -
wget https://archive.apache.org/dist/spark/spark-4.1.2/spark-4.1.2-bin-hadoop3.tgz
tar -xvzf spark-4.1.2-bin-hadoop3.tgz
先决条件 CDM 5.x+
- Java11(最低要求),因为 Spark 3.x 二进制文件是使用它编译的。
- 如有需要,可以使用任何 Java 11+ LTS 版本(例如 Java 17、Java 21 或 Java 25)
- Spark 版本 3.5.8 搭配 Scala 2.13
可以通过运行以下命令安装 Spark: -
wget https://archive.apache.org/dist/spark/spark-3.5.8/spark-3.5.8-bin-hadoop3-scala2.13.tgz
tar -xvzf spark-3.5.8-bin-hadoop3-scala2.13.tgz
[!CAUTION] 如果存在任何版本不匹配,您可能会看到如下异常,
Exception in thread "main" java.lang.NoSuchMethodError: 'void scala.runtime.Statics.releaseFence()'
注意:
- 通常使用上述二进制文件在单个 VM(无需集群)上安装,该 VM 是您希望运行此作业的位置。对于大多数一次性迁移,推荐此简单配置。
- 对于大型(例如数 TB)复杂迁移,或者当 CDM 用作长期数据传输工具而非一次性作业时,您可以使用支持上述版本的 Spark 集群或 Spark Serverless 平台,如
Watsonx.data、Databricks或Google Dataproc。 - 在 Spark 集群上部署 CDM 时,将参数
--master "local[*]"替换为--master "spark://master-host:port",并移除与单 VM 运行相关的任何参数(例如--driver-memory、--executor-memory等)
数据迁移步骤:
cdm.propertiesfile needs to be configured as applicable for the environment. The file can have any name, it does not need to becdm.properties.- A sample properties file with default values can be found here as cdm.properties
- A complete reference properties file with default values can be found here as cdm-detailed.properties
- Place the properties file where it can be accessed while running the job via spark-submit.
- Run the job using
spark-submitcommand as shown below:
spark-submit --properties-file cdm.properties \
--conf spark.cdm.schema.origin.keyspaceTable="<keyspacename>.<tablename>" \
--master "local[*]" --driver-memory 25G --executor-memory 25G \
--class com.datastax.cdm.job.Migrate cassandra-data-migrator-5.x.x.jar &> logfile_name_$(date +%Y%m%d_%H_%M).txt
注意:
- 上述命令会生成日志文件
logfile_name_*.txt,以避免在控制台输出日志。 - 根据您的使用场景更新内存选项(驱动程序 & 执行器内存)
- 若要跟踪某次运行的详细信息(记录在
target键空间中),请传递参数--conf spark.cdm.trackRun=true - 若要仅筛选特定令牌范围内的记录,请向
Migration或Validation作业传递以下两个额外参数
--conf spark.cdm.filter.cassandra.partition.min=<token-range-min>
--conf spark.cdm.filter.cassandra.partition.max=<token-range-max>
数据验证的步骤:
- 要在数据验证模式下运行作业,请使用类选项
--class com.datastax.cdm.job.DiffData,如下所示
spark-submit --properties-file cdm.properties \
--conf spark.cdm.schema.origin.keyspaceTable="<keyspacename>.<tablename>" \
--master "local[*]" --driver-memory 25G --executor-memory 25G \
--class com.datastax.cdm.job.DiffData cassandra-data-migrator-5.x.x.jar &> logfile_name_$(date +%Y%m%d_%H_%M).txt
- 验证作业将在日志文件中将差异报告为“ERRORS”,如下所示。
23/04/06 08:43:06 ERROR DiffJobSession: Mismatch row found for key: [key3] Mismatch: Target Index: 1 Origin: valueC Target: value999)
23/04/06 08:43:06 ERROR DiffJobSession: Corrected mismatch row in target: [key3]
23/04/06 08:43:06 ERROR DiffJobSession: Missing target row found for key: [key2]
23/04/06 08:43:06 ERROR DiffJobSession: Inserted missing row in target: [key2]
- 请从输出日志文件中 grep 所有
ERROR,以获取缺失和不匹配记录的列表。- 请注意,它按主键值列出差异。
- 如果您希望将此类日志(包含
missing和mismatched行详情的行)重定向到单独的文件,您可以使用log4j2.properties文件 在此提供 如下所示
spark-submit --properties-file cdm.properties \
--conf spark.cdm.schema.origin.keyspaceTable="<keyspacename>.<tablename>" \
--conf spark.executor.extraJavaOptions='-Dlog4j2.configurationFile=log4j2.properties' \
--conf spark.driver.extraJavaOptions='-Dlog4j2.configurationFile=log4j2.properties' \
--master "local[*]" --driver-memory 25G --executor-memory 25G \
--class com.datastax.cdm.job.DiffData cassandra-data-migrator-5.x.x.jar &> logfile_name_$(date +%Y%m%d_%H_%M).txt
- Validation 作业也可以在 AutoCorrect 模式下运行。此模式可以
- 将
origin中缺失的任何记录添加到target - 更新
origin和target之间不匹配的任何记录
- 将
- 使用 properties 文件中的以下一个或两个参数启用/禁用此功能
spark.cdm.autocorrect.missing false|true
spark.cdm.autocorrect.mismatch false|true
[!IMPORTANT] Validation 作业永远不会从目标中删除记录,即它仅在目标上添加或更新数据
重新运行(先前未完成的)迁移或验证
- 您可以通过将
spark.cdm.trackRun.autoRerun设置为true(默认为false)来重新运行/恢复因任何原因停止(或带有一些错误完成)的上一个迁移或验证作业,如下所示。这将自动发现上一个作业的进度,并从该点恢复,即它会跳过上一次运行中已成功迁移(或验证)的任何 token-ranges。
spark-submit --properties-file cdm.properties \
--conf spark.cdm.schema.origin.keyspaceTable="<keyspacename>.<tablename>" \
--conf spark.cdm.trackRun.autoRerun=true \
--master "local[*]" --driver-memory 25G --executor-memory 25G \
--class com.datastax.cdm.job.<Migrate|DiffData> cassandra-data-migrator-5.x.x.jar &> logfile_name_$(date +%Y%m%d_%H_%M).txt
- 你还可以通过传递特定的
spark.cdm.trackRun.previousRunId来恢复某个特定的先前作业(而非最后一个作业),该作业可能因任何原因而停止(或带有一些错误完成),如下所示
spark-submit --properties-file cdm.properties \
--conf spark.cdm.schema.origin.keyspaceTable="<keyspacename>.<tablename>" \
--conf spark.cdm.trackRun.previousRunId=<prev_run_id> \
--master "local[*]" --driver-memory 25G --executor-memory 25G \
--class com.datastax.cdm.job.<Migrate|DiffData> cassandra-data-migrator-5.x.x.jar &> logfile_name_$(date +%Y%m%d_%H_%M).txt
执行大字段 Guardrail 违规检查
- 此模式可帮助识别
origin表中可能破坏集群 Guardrail 的大字段(例如 AstraDB 对单个大字段有 10MB 限制),如下所示使用类选项--class com.datastax.cdm.job.GuardrailCheck
spark-submit --properties-file cdm.properties \
--conf spark.cdm.schema.origin.keyspaceTable="<keyspacename>.<tablename>" \
--conf spark.cdm.feature.guardrail.colSizeInKB=10000 \
--master "local[*]" --driver-memory 25G --executor-memory 25G \
--class com.datastax.cdm.job.GuardrailCheck cassandra-data-migrator-5.x.x.jar &> logfile_name_$(date +%Y%m%d_%H_%M).txt
[!NOTE] 此模式仅操作一个数据库,即
origin,此模式下不存在target
功能
- 自动检测表架构(列名、类型、键、集合、UDT 等)
- 包括计数器表 Counter tables
- 重新运行/恢复之前可能因任何原因(被终止、出现异常等)而停止的作业
- 如果重新运行
validation作业,它将包含在上一次运行中存在差异的任何 token 范围 - 如果重新运行反复失败,您可以使用
rerunMultiplier功能来帮助提高成功几率。
- 如果重新运行
- 保留 writetimes 和 TTLs
- 支持高级 DataTypes 的迁移/验证(Sets、Lists、Maps、UDTs)
- 使用
writetime和/或 CQL 条件以及/或 token 范围列表从Origin中过滤记录 - 执行护栏检查(识别大字段)
- 支持在
Target上添加constants作为新列 - 支持将
Origin上的Map列扩展为Target上的多条记录 - 支持从
Origin中的 JSON 列提取值并将其映射到Target上的特定字段 - 可部署在 Spark 集群或单台虚拟机上
- 完全容器化(兼容 Docker 和 K8s)
- SSL 支持(包括自定义加密算法)
- 从任何 Cassandra
Origin(Apache Cassandra® / DataStax Enterprise™ / DataStax Astra DB™)迁移到任何 CassandraTarget(Apache Cassandra® / DataStax Enterprise™ / DataStax Astra DB™) - 使用 DevOps API 自动下载 Astra DB 的 Secure Connect Bundles
- 支持从 Azure Cosmos Cassandra 进行迁移/验证
- 使用较小的随机化数据集验证迁移的准确性和性能
- 支持添加自定义固定
writetime和/或ttl - 在目标 keyspace 的表(
cdm_run_info和cdm_run_details)中跟踪运行信息(开始时间、结束时间、运行指标、状态等)
须知
[!TIP] 如果您想传入额外的 Cassandra Java Driver 配置,您可以传入
--conf spark.driver.extraJavaOptions="-Ddatastax-java-driver.advanced.connection.pool.remote.size=5
[!TIP] 如果您想记录
DEBUG级别的语句以进行故障排除,您可以传入--conf spark.driver.extraJavaOptions='-Dlog4j2.level=DEBUG -Dlog4j2.rootLogger.level=DEBUG' --conf spark.executor.extraJavaOptions='-Dlog4j2.level=DEBUG -Dlog4j2.rootLogger.level=DEBUG'
- 每次运行(迁移或验证)均可被跟踪(当启用时)。您可以在目标键空间中的表
cdm_run_info和cdm_run_details中找到其摘要和详细信息。 - 出于优化原因,CDM 不会在字段级别迁移
ttl和writetime。相反,它会找到origin行中最高ttl的字段和最高writetime的字段,并将这些值应用于整个target行。 - 出于性能原因,CDM 默认忽略使用集合和 UDT 字段进行
ttl和writetime计算。如果您希望包含此类字段,请将spark.cdm.schema.ttlwritetime.calc.useCollections参数设置为true。 - 如果表仅包含集合和/或 UDT 非键列,且没有表级
ttl配置,则目标将没有ttl,这可能导致origin和target之间出现不一致,因为行会因ttl过期而在origin上过期。如果您希望避免这种情况,我们建议在此类场景中将spark.cdm.schema.ttlwritetime.calc.useCollections参数设置为true。 - 如果表仅包含集合和/或 UDT 非键列,则在目标上使用的
writetime将是作业运行的时间。如果您希望避免这种情况,我们建议在此类场景中将spark.cdm.schema.ttlwritetime.calc.useCollections参数设置为true。 - CDM 使用
UNSET值作为空字段(包括空文本)的值,以避免在创建行时创建(或延续)墓碑。 - 当对同一张表多次运行 CDM 迁移(或带自动更正的验证)时(无论出于何种原因),可能会导致
list类型列中出现重复条目。请注意,这是 由 Cassandra/DSE 的 bug 引起,而非 CDM 的问题。可以通过启用并设置spark.cdm.transform.custom.writetime.incrementBy参数为正值来解决此问题。该参数是专门为此问题而添加的。 - 当您重新运行作业以从上一次运行恢复时,在表
cdm_run_info中捕获的运行指标(读取、写入、跳过等)仅针对当前运行。如果上一次运行因某些原因被终止,其运行指标可能未被保存。如果上一次运行已完成(未被终止)但存在错误,则您还将拥有上一次运行的所有运行指标。 - 在 Spark 集群上运行时(而非单台虚拟机),速率限制值(
spark.cdm.perfops.ratelimit.origin和spark.cdm.perfops.ratelimit.target)适用于各个 Spark 工作节点。因此,该值应设置为 所需的有效速率限制/Spark 工作节点数量。例如,如果您需要 10000 的有效速率限制,且 Spark 工作节点数量为 4,则应将上述速率限制参数设置为 2500。
性能建议
以下建议仅在迁移大型表且默认性能不足时可能有用
- 性能瓶颈通常是由以下原因造成的
Origin或Target集群上的资源可用性低- CDM 虚拟机上的资源可用性低,参见此处建议
- 糟糕的模式设计,这可能是由失衡的
Origin集群、大分区(> 100 MB)、大行(> 10MB)和/或高列数引起的。
- 以下属性的错误配置可能会对性能产生负面影响
numParts: 默认值为 5K,但理想值通常约为 table-size/10MB。batchSize: 默认值为 5,但对于 primary-key=partition-key 的表或平均行大小 > 20 KB 的表,应将其设置为 1。同样,如果行大小较小(< 1KB)且大多数分区包含多行(100+),则应将其设置为大于 5 的值。fetchSizeInRows: 默认值为 1K,这通常运行良好。但是,如果您的表有许多大行(超过 100KB),您可以按需减少此值。ratelimit: 默认值为20000,但通常应在更新其他属性后,将此属性更新为origin和target集群能够高效处理的最大可能值。
- 使用模式操作功能(如
constantColumns、explodeMap、extractJson)、转换函数和/或 where 过滤条件(分区最小值/最大值除外)可能会对性能产生负面影响 - 我们通常建议为 CDM VM 使用 此基础设施 和 此初始配置。然后,您可以基于上述提供的 CDM 参数信息以及
Origin和Target集群上观察到的负载和吞吐量,进一步优化作业 - 对于大型(例如数 TB)复杂迁移,或者当 CDM 被用作长期数据传输工具而非一次性作业时,我们建议使用 Spark 集群或 Spark Serverless 平台,如
Watsonx.data或Databricks或Google Dataproc。
[!NOTE] 如需进行额外的性能调优,请参阅此处
cdm-detailed.properties文件 中提到的详细信息
Building Jar for local development
- Clone this repo
- Move to the repo folder
cd cassandra-data-migrator - Run the build
mvn clean package(Needs Maven 3.9.x) - The fat jar (
cassandra-data-migrator-5.x.x.jar) file should now be present in thetargetfolder
贡献者
在此](./CONTRIBUTING.md#contributors)查看我们所有出色的贡献者。