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") )
方案说明
- 两种方案均无需重复读取Kinesis流,基于同一微批次数据流处理;
- 方案1代码简洁、兼容性好,适合需要标准SQL聚合、允许Shuffle的场景;
- 方案2利用Kinesis与Spark分区的对应关系,避免Shuffle,性能更优,适合纯内存级轻量计算;
- 若需同时保留明细与聚合数据,启动两个独立写流任务即可,Spark会自动优化输入复用。
内容的提问来源于stack exchange,提问作者Mark J Miller
相关产品推荐
相关产品推荐

