You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Delta Live Tables工作流中将Gold表数据同步至Cassandra?

在Delta Live Tables工作流中将Gold表数据同步到Cassandra

我们现有一套Delta Live Tables(DLT)工作流,从Kafka Topic读取数据后,经过清洗、过滤、处理、聚合,最终写入Bronze/Silver/Gold层表。现在需要把Gold表的聚合结果导出到Cassandra表,用于构建REST查询服务。我尝试修改Gold表脚本,在写入Gold表后添加同步到Cassandra的步骤,但没有成功,尝试的代码如下:

@dlt.table
def test_live_gold():
  return (
    dlt.read("test_kafka_silver").groupBy("user_id", "event_type").count()

#     df = spark.read.format("delta")
#       .table("customer.test_live_gold")
#       .withColumnRenamed("user_id", "account_id")
#       .withColumnRenamed("event_type", "event_name")
#       .withColumn("last_updated_dt", current_timestamp())
#     df.show(5, False)
#     write_to_cassandra_table('customer', 'test_keyspace', df)
  )

请问如何在同一个DLT工作流里实现从Delta表到Cassandra的数据复制?


方法1:使用DLT run_after 装饰器执行独立同步任务

DLT提供的run_after装饰器可以指定在目标表完成构建后,触发后续任务,这种方式能保证Gold表的构建与Cassandra同步逻辑隔离,避免同步失败影响Gold表生成:

@dlt.table
def test_live_gold():
  return (
    dlt.read("test_kafka_silver")
      .groupBy("user_id", "event_type")
      .count()
  )

@dlt.run_after("test_live_gold")
def sync_gold_to_cassandra():
  # 读取刚生成的Gold表数据并做字段转换
  sync_df = dlt.read("test_live_gold") \
    .withColumnRenamed("user_id", "account_id") \
    .withColumnRenamed("event_type", "event_name") \
    .withColumn("last_updated_dt", current_timestamp())
  
  # 写入Cassandra,若有自定义write函数可直接替换此处逻辑
  sync_df.write \
    .format("org.apache.spark.sql.cassandra") \
    .mode("append")  # 按需选择append/overwrite/upsert
    .options(
      table="your_cassandra_target_table",
      keyspace="test_keyspace",
      spark_cassandra_connection_host="your_cassandra_host_ip"
    ) \
    .save()

方法2:在Gold表定义中嵌入同步逻辑(不推荐)

如果希望在生成Gold表的同时完成同步,可以在表定义中先生成目标DataFrame,再执行同步操作,最后返回用于构建Gold表的DataFrame。但这种方式会让同步失败直接导致Gold表构建失败,不符合分层隔离的设计原则:

@dlt.table
def test_live_gold():
  # 生成Gold表的核心数据
  gold_df = dlt.read("test_kafka_silver") \
    .groupBy("user_id", "event_type") \
    .count()
  
  # 同步到Cassandra的字段转换与写入
  sync_df = gold_df \
    .withColumnRenamed("user_id", "account_id") \
    .withColumnRenamed("event_type", "event_name") \
    .withColumn("last_updated_dt", current_timestamp())
  
  sync_df.write \
    .format("org.apache.spark.sql.cassandra") \
    .mode("append")
    .options(
      table="your_cassandra_target_table",
      keyspace="test_keyspace",
      spark_cassandra_connection_host="your_cassandra_host_ip"
    ) \
    .save()
  
  # 返回用于构建Gold表的DataFrame
  return gold_df

关键注意事项

  • 确保Databricks集群已安装对应版本的Cassandra Spark连接器,例如Maven坐标com.datastax.spark:spark-cassandra-connector_2.12:3.4.1(需匹配Spark版本)。
  • 写入模式选择:append适合增量同步,overwrite适合全量覆盖,若需upsert需结合Cassandra主键配置对应参数。
  • 大规模数据场景下,建议添加过滤逻辑(如基于时间戳增量同步),避免每次同步全量数据。

内容的提问来源于stack exchange,提问作者user468587

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.20 00:30:57