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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 08:50:37