如何用PySpark每15分钟读取S3中IoT数据并计算均值
解决方案
核心思路
要实现每15分钟处理对应时段的新增数据,关键要做到两点:精准过滤目标时段的数据和避免重复处理已计算过的数据。结合你的场景,分以下几步实现:
1. 精准过滤目标时段数据
如果你的Kinesis Firehose已经按时间分区写入S3(比如year=YYYY/month=MM/day=DD/hour=HH/minute=MM这种细粒度分区),可以直接通过S3路径过滤,大幅减少数据读取量;如果没有配置分区,就通过数据里的time字段进行过滤。
方式一:利用S3分区过滤(推荐)
假设Firehose配置了15分钟粒度的分区(若未配置,建议调整Firehose的分区设置提升效率),动态生成目标时段的S3路径:
import pyspark from pyspark.sql import SparkSession from pyspark.sql import functions as f from datetime import datetime, timedelta # 实际调度时可从调度工具(如Airflow)或环境变量获取当前时段 target_end_time = datetime(2023, 1, 17, 0, 15) # 当前调度时间,对应处理00:00-00:15的时段 target_start_time = target_end_time - timedelta(minutes=15) # 生成匹配目标时段的S3分区路径(根据你的实际分区格式调整) s3_path = f"s3a://<my-bucket-name>/<sitename>/year={target_start_time.year}/month={target_start_time.month:02d}/day={target_start_time.day:02d}/hour={target_start_time.hour:02d}/minute={target_start_time.minute:02d}/*.parquet" spark = SparkSession.builder.appName('iot_15min_avg').getOrCreate() # 只读取目标时段的分区数据 df = spark.read.parquet(s3_path) df = df.withColumn("time", f.to_timestamp("time", 'dd-MM-yyyy HH:mm:ss'))
方式二:通过time字段过滤(无分区时用)
如果没有S3分区,先读取当日数据再过滤目标时段:
import pyspark from pyspark.sql import SparkSession from pyspark.sql import functions as f from datetime import datetime, timedelta target_end_time = datetime(2023, 1, 17, 0, 15) target_start_time = target_end_time - timedelta(minutes=15) # 读取当日所有数据 s3_path = f"s3a://<my-bucket-name>/<sitename>/year={target_start_time.year}/month={target_start_time.month:02d}/day={target_start_time.day:02d}/*.parquet" spark = SparkSession.builder.appName('iot_15min_avg').getOrCreate() df = spark.read.parquet(s3_path) df = df.withColumn("time", f.to_timestamp("time", 'dd-MM-yyyy HH:mm:ss')) # 过滤出目标15分钟时段的数据 filtered_df = df.filter( (f.col("time") >= target_start_time) & (f.col("time") < target_end_time) )
2. 计算15分钟均值
过滤出目标时段数据后,批量计算所有参数的均值:
# 提取除time外的所有参数字段 metric_cols = [col for col in filtered_df.columns if col != "time"] # 计算每个参数的均值并命名 avg_df = filtered_df.agg(*[f.avg(col).alias(f"{col}_avg") for col in metric_cols]) # 添加时段标识字段,方便后续查询 avg_df = avg_df.withColumn("start_time", f.lit(target_start_time.strftime('%Y-%m-%d %H:%M:%S'))) avg_df = avg_df.withColumn("end_time", f.lit(target_end_time.strftime('%Y-%m-%d %H:%M:%S')))
3. 写入结果到目标S3桶
将均值结果写入另一个S3桶,建议按时间分区存储以便后续查询:
target_s3_path = f"s3a://<target-bucket-name>/iot_avg/year={target_start_time.year}/month={target_start_time.month:02d}/day={target_start_time.day:02d}/" # 用append模式写入,避免覆盖已有数据 avg_df.write.mode("append").parquet(target_s3_path)
4. 避免重复处理
为防止同一时段数据被重复计算,需跟踪已处理的时段:
- 在目标S3桶的根目录下维护一个
processed_intervals.txt文件,记录已处理的start_time-end_time字符串 - 每次调度前先读取该文件,判断当前目标时段是否已处理,若已处理则跳过
- 处理完成后,将当前时段写入该文件
示例代码片段:
import boto3 s3 = boto3.client('s3') processed_file_key = "iot_avg/processed_intervals.txt" current_interval = f"{target_start_time.strftime('%Y-%m-%d %H:%M:%S')}-{target_end_time.strftime('%Y-%m-%d %H:%M:%S')}" # 读取已处理时段列表 try: response = s3.get_object(Bucket='<target-bucket-name>', Key=processed_file_key) processed_intervals = response['Body'].read().decode('utf-8').splitlines() except s3.exceptions.NoSuchKey: processed_intervals = [] if current_interval in processed_intervals: print(f"时段 {current_interval} 已处理,跳过") spark.stop() exit() # 执行数据处理和写入逻辑... # 记录当前已处理的时段 processed_intervals.append(current_interval) s3.put_object( Bucket='<target-bucket-name>', Key=processed_file_key, Body='\n'.join(processed_intervals) )
调度建议
可以用Airflow、AWS EventBridge等调度工具,每15分钟触发一次PySpark任务,任务执行时自动计算当前要处理的15分钟时段(比如当前时间00:15,处理00:00-00:15;当前时间00:30,处理00:15-00:30)。
内容的提问来源于stack exchange,提问作者anonymous_33008899
相关产品推荐
相关产品推荐

