ITADN
datastax/cassandra-data-migrator
datastax/cassandra-data-migrator · 文件 下载 ZIP
文件最后提交记录最后更新时间
README.md
以下内容由 AI 翻译,如有问题请点此提交 issue 反馈

License Apache2 Star on Github GitHub release (with filter) Docker Pulls

cassandra-data-migrator (also known as CDM)

在源和目标 Cassandra 集群之间迁移和验证表。

[!IMPORTANT] 请注意以下列出的先决条件,以避免 CDM、Spark、Scala、Hadoop 版本不匹配问题

作为容器安装

  • DockerHub 获取包含所有依赖项的最新镜像
    • 所有迁移工具(cassandra-data-migrator + dsbulk + cqlsh)将在容器的 /assets/ 文件夹中可用

作为 JAR 文件安装

  • 使用以下方法之一下载最新的 jar 文件 GitHub release (with filter)
    • 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.dataDatabricksGoogle Dataproc
  • 在 Spark 集群上部署 CDM 时,将参数 --master "local[*]" 替换为 --master "spark://master-host:port",并移除与单 VM 运行相关的任何参数(例如 --driver-memory--executor-memory 等)

数据迁移步骤:

  1. cdm.properties file needs to be configured as applicable for the environment. The file can have any name, it does not need to be cdm.properties.
  2. Place the properties file where it can be accessed while running the job via spark-submit.
  3. Run the job using spark-submit command 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
  • 若要仅筛选特定令牌范围内的记录,请向 MigrationValidation 作业传递以下两个额外参数
--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,以获取缺失和不匹配记录的列表。
    • 请注意,它按主键值列出差异。
  • 如果您希望将此类日志(包含 missingmismatched 行详情的行)重定向到单独的文件,您可以使用 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
    • 更新 origintarget 之间不匹配的任何记录
  • 使用 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 等)
  • 重新运行/恢复之前可能因任何原因(被终止、出现异常等)而停止的作业
    • 如果重新运行 validation 作业,它将包含在上一次运行中存在差异的任何 token 范围
    • 如果重新运行反复失败,您可以使用 rerunMultiplier 功能来帮助提高成功几率。
  • 保留 writetimesTTLs
  • 支持高级 DataTypes 的迁移/验证(SetsListsMapsUDTs
  • 使用 writetime 和/或 CQL 条件以及/或 token 范围列表从 Origin 中过滤记录
  • 执行护栏检查(识别大字段)
  • 支持在 Target 上添加 constants 作为新列
  • 支持将 Origin 上的 Map 列扩展为 Target 上的多条记录
  • 支持从 Origin 中的 JSON 列提取值并将其映射到 Target 上的特定字段
  • 可部署在 Spark 集群或单台虚拟机上
  • 完全容器化(兼容 Docker 和 K8s)
  • SSL 支持(包括自定义加密算法)
  • 从任何 Cassandra OriginApache Cassandra® / DataStax Enterprise™ / DataStax Astra DB™)迁移到任何 Cassandra TargetApache Cassandra® / DataStax Enterprise™ / DataStax Astra DB™
  • 使用 DevOps API 自动下载 Astra DB 的 Secure Connect Bundles
  • 支持从 Azure Cosmos Cassandra 进行迁移/验证
  • 使用较小的随机化数据集验证迁移的准确性和性能
  • 支持添加自定义固定 writetime 和/或 ttl
  • 在目标 keyspace 的表(cdm_run_infocdm_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_infocdm_run_details 中找到其摘要和详细信息。
  • 出于优化原因,CDM 不会在字段级别迁移 ttlwritetime。相反,它会找到 origin 行中最高 ttl 的字段和最高 writetime 的字段,并将这些值应用于整个 target 行。
  • 出于性能原因,CDM 默认忽略使用集合和 UDT 字段进行 ttlwritetime 计算。如果您希望包含此类字段,请将 spark.cdm.schema.ttlwritetime.calc.useCollections 参数设置为 true
  • 如果表仅包含集合和/或 UDT 非键列,且没有表级 ttl 配置,则目标将没有 ttl,这可能导致 origintarget 之间出现不一致,因为行会因 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.originspark.cdm.perfops.ratelimit.target)适用于各个 Spark 工作节点。因此,该值应设置为 所需的有效速率限制/Spark 工作节点数量。例如,如果您需要 10000 的有效速率限制,且 Spark 工作节点数量为 4,则应将上述速率限制参数设置为 2500。

性能建议

以下建议仅在迁移大型表且默认性能不足时可能有用

  • 性能瓶颈通常是由以下原因造成的
    • OriginTarget 集群上的资源可用性低
    • 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,但通常应在更新其他属性后,将此属性更新为 origintarget 集群能够高效处理的最大可能值。
  • 使用模式操作功能(如 constantColumnsexplodeMapextractJson)、转换函数和/或 where 过滤条件(分区最小值/最大值除外)可能会对性能产生负面影响
  • 我们通常建议为 CDM VM 使用 此基础设施此初始配置。然后,您可以基于上述提供的 CDM 参数信息以及 OriginTarget 集群上观察到的负载和吞吐量,进一步优化作业
  • 对于大型(例如数 TB)复杂迁移,或者当 CDM 被用作长期数据传输工具而非一次性作业时,我们建议使用 Spark 集群或 Spark Serverless 平台,如 Watsonx.dataDatabricksGoogle Dataproc

[!NOTE] 如需进行额外的性能调优,请参阅此处 cdm-detailed.properties 文件 中提到的详细信息

Building Jar for local development

  1. Clone this repo
  2. Move to the repo folder cd cassandra-data-migrator
  3. Run the build mvn clean package (Needs Maven 3.9.x)
  4. The fat jar (cassandra-data-migrator-5.x.x.jar) file should now be present in the target folder

贡献者

在此](./CONTRIBUTING.md#contributors)查看我们所有出色的贡献者。