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

Spark Streaming单遍处理:写入源记录并按Kinesis分片计算聚合

现有流处理管道实现

读取Kinesis流代码

import pyspark.sql.types as T
import pyspark.sql.functions as F

schema = T.StructType(
    [
        T.StructField("data", T.StringType(), True),
        T.StructField("metadata", T.StructType([
            T.StructField("timestamp", T.TimestampType(), True),
            T.StructField("record-type", T.StringType(), True),
            T.StructField("operation", T.StringType(), True),
            T.StructField("partition-key-type", T.StringType(), True),
            T.StructField("partition-key-value", T.StringType(), True),
            T.StructField("schema-name", T.StringType(), True),
            T.StructField("table-name", T.StringType(), True)
        ]), True)
    ]
)

kinesis = (
    spark.readStream
        .format("kinesis")
        .option("streamName", "data-platform-aws-dms-poc-dms-sample")
        .option("region", "us-east-1")
        .option("initialPosition", "TRIM_HORIZON")
        .load()
        .selectExpr("partitionKey", "CAST(data AS STRING) AS data", "stream", "shardId", "sequenceNumber", "approximateArrivalTimestamp")
        .withColumn("data", F.from_json("data", schema))
        .selectExpr("partitionKey", "data.data", "data.metadata.*", "stream", "shardId", "sequenceNumber", "approximateArrivalTimestamp")
)

写入Delta Lake代码

checkpoint_path = "s3://route-databricks-sandbox-data/checkpoint/sandbox.route.dms_sample"

query = (
    kinesis
        .writeStream
        .format("delta")
        .queryName("dms_sample_output")
        .outputMode("append")
        .option("checkpointLocation", checkpoint_path)
        .partitionBy("schema-name", "table-name")
        .toTable("sandbox.public.dms_sample")
)
需求

在单次流处理过程中,无需重新读取Kinesis流数据,针对每个Kinesis shardId执行数据量较小、可轻松放入内存的聚合计算。

可行实现方案

方案1:Spark流式分组聚合

直接在读取后的DataFrame上按shardId分组执行轻量聚合,可选择将聚合结果与原明细数据结合写入,或单独输出聚合结果。Spark会在同一个微批次内处理数据,不会重复读取Kinesis流。

示例代码(以统计shard内记录数、操作类型分布为例):

# 计算shard级聚合结果
agg_df = kinesis.groupBy("shardId") \
    .agg(
        F.count("*").alias("shard_record_count"),
        F.collect_set("operation").alias("operation_types"),
        F.min("approximateArrivalTimestamp").alias("shard_min_arrival_time"),
        F.max("approximateArrivalTimestamp").alias("shard_max_arrival_time")
    )

# 方式A:聚合结果与原明细关联后写入(保留明细+聚合信息)
combined_df = kinesis.join(agg_df, on="shardId", how="left")
combined_query = (
    combined_df
        .writeStream
        .format("delta")
        .queryName("dms_sample_with_shard_agg")
        .outputMode("append")
        .option("checkpointLocation", checkpoint_path + "_with_agg")
        .partitionBy("schema-name", "table-name")
        .toTable("sandbox.public.dms_sample_with_shard_agg")
)

# 方式B:单独输出聚合结果(仅需聚合数据场景)
agg_query = (
    agg_df
        .writeStream
        .format("delta")
        .queryName("shard_agg_output")
        .outputMode("complete")  # 分组聚合用complete模式输出全量聚合结果
        .option("checkpointLocation", checkpoint_path + "_shard_agg")
        .toTable("sandbox.public.shard_agg_results")
)

方案2:基于Spark分区的本地聚合

Spark读取Kinesis时,每个Kinesis Shard对应一个Spark RDD Partition,可利用mapPartitions算子在分区内直接做内存级聚合,避免Shuffle操作,性能更优。

示例代码:

from pyspark.sql import Row

def shard_agg(iterator):
    shard_id = None
    record_count = 0
    operation_counts = {}
    min_arrival_time = None
    max_arrival_time = None
    
    for row in iterator:
        if shard_id is None:
            shard_id = row.shardId
        
        record_count += 1
        op = row.operation
        operation_counts[op] = operation_counts.get(op, 0) + 1
        
        arrival_time = row.approximateArrivalTimestamp
        if min_arrival_time is None or arrival_time < min_arrival_time:
            min_arrival_time = arrival_time
        if max_arrival_time is None or arrival_time > max_arrival_time:
            max_arrival_time = arrival_time
    
    yield Row(
        shardId=shard_id,
        record_count=record_count,
        operation_counts=str(operation_counts),
        min_arrival_time=min_arrival_time,
        max_arrival_time=max_arrival_time
    )

# 转换为RDD执行分区聚合,再转回DataFrame
shard_agg_df = spark.createDataFrame(kinesis.rdd.mapPartitions(shard_agg))

# 写入聚合结果
agg_query = (
    shard_agg_df
        .writeStream
        .format("delta")
        .queryName("shard_partition_agg_output")
        .outputMode("append")  # 输出当前微批次的shard聚合结果
        .option("checkpointLocation", checkpoint_path + "_partition_agg")
        .toTable("sandbox.public.shard_partition_agg_results")
)

# 同时保留原明细数据,启动独立写流任务即可(共享同一输入流,无重复读取)
original_query = (
    kinesis
        .writeStream
        .format("delta")
        .queryName("dms_sample_output")
        .outputMode("append")
        .option("checkpointLocation", checkpoint_path)
        .partitionBy("schema-name", "table-name")
        .toTable("sandbox.public.dms_sample")
)

方案说明

  1. 两种方案均无需重复读取Kinesis流,基于同一微批次数据流处理;
  2. 方案1代码简洁、兼容性好,适合需要标准SQL聚合、允许Shuffle的场景;
  3. 方案2利用Kinesis与Spark分区的对应关系,避免Shuffle,性能更优,适合纯内存级轻量计算;
  4. 若需同时保留明细与聚合数据,启动两个独立写流任务即可,Spark会自动优化输入复用。

内容的提问来源于stack exchange,提问作者Mark J Miller

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 23:40:19