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

GCP Dataproc上Spark Structured Streaming触发间隔不生效问题排查

问题:Spark Structured Streaming触发间隔30秒实际却90秒执行一次

环境背景

在GCP Compute Engine的单节点(1主0从)Dataproc集群上,运行一个从Pub/Sub Lite读取数据的Structured Streaming任务,核心逻辑实现于transform.py。

核心代码简化版

class Transform:
    def __init__(self, app_name):
        self.app_name = app_name

    def create_spark_session(self):
        spark_session = SparkSession\
            .builder\
            .appName(self.app_name)\
            .master("yarn")\
            .getOrCreate()
        return spark_session

    @staticmethod
    def read_stream_from_pubsub(spark_session, project_number, region, subscription) -> DataFrame:
        df = (
            spark_session
            .readStream
            .format("pubsublite")
            .option(
                "pubsublite.subscription",
                f"projects/{project_number}/locations/{region}/subscriptions/{subscription}",
            )
            .load()
        )
        return df

    @staticmethod
    def process_df(df) -> DataFrame:
        ticker_schema = Schemas.ticker_schema()
        ticker_df = df \
            .select(
                f.from_json(f.col('data').cast('string'), ticker_schema).alias('data')
            )
        return ticker_df

    @staticmethod
    def process_df_cast(df_process) -> DataFrame:
        df_process_cast = df_process \
            .select(
                'data.*'
            ) \
            .withColumn('time', f.to_timestamp(f.col('time'))) \
            .withColumn('string_real_time', f.col('time')[0:19]) \
            .withColumn('price_now', f.col('price_now').cast('double')) \
            .withColumn('open', f.col('open').cast('double')) \
            .withColumn('close', f.col('close').cast('double'))

        df_process_cast_group_by_window = df_process_cast \
            .groupBy(
                f.window(f.col('time'), '30 seconds').alias('ventana'),
                f.col('localSymbol'),
                f.col('close'),
                f.col('open')
            ) \
            .agg(
                f.max(f.col('price_now')).alias('max_price_window'),
                f.min(f.col('price_now')).alias('min_price_window'),
                f.first(f.col('price_now')).alias('start_price_window'),
                f.last(f.col('price_now')).alias('end_price_window'),
            )

        df_process_cast_group_by_window = df_process_cast_group_by_window\
            .select(
                f.col('*')
            ) \
            .withColumn('start_window', f.col('ventana.start')) \
            .withColumn('end_window', f.col('ventana.end'))

        return df_process_cast_group_by_window

    def main(self):
        ss = self.create_spark_session()
        df = self.read_stream_from_pubsub(ss, settings.PROJECT_NUMBER, settings.REGION, settings.SUBSCRIPTION)
        ticker_df = self.process_df(df)
        ticker_df_cast = self.process_df_cast(ticker_df)

        query = ticker_df_cast.writeStream \
            .foreachBatch(self.analysis_batch) \
            .outputMode('update') \
            .trigger(processingTime='30 seconds') \
            .start()
        query.awaitTermination()

    def analysis_batch(self, ticker_df_cast, epoch_id):
        print("im here")

if __name__ == '__main__':
    transform = Transform("APP_IB_SCALPING")
    transform.main()

问题现象

代码中明确设置trigger(processingTime='30 seconds'),但任务实际每隔90秒才执行一次,为设定间隔的三倍。

已完成的排查动作

  • 反复核对代码,确认触发间隔配置无误;
  • 检查Dataproc集群指标,未发现资源占用过高导致的延迟;
  • 验证Pub/Sub Topic消息发布流程,确认无延迟;
  • 查看Spark日志,未找到与定时逻辑相关的报错或警告。

预期效果

任务严格按照30秒的触发间隔执行。

任务提交代码

...
@staticmethod
def submit_job(project_id, region, cluster_name, gcs_bucket, spark_filename):
    # 创建作业客户端
    job_client = dataproc_v1.JobControllerClient(
        client_options={"api_endpoint": "{}-dataproc.googleapis.com:443".format(region)}
    )

    # 配置作业参数
    job = {
        "placement": {"cluster_name": cluster_name},
        "pyspark_job": {
            "main_python_file_uri": "gs://{}/{}".format(gcs_bucket, spark_filename),
            "jar_file_uris": [
                "gs://spark-lib/pubsublite/pubsublite-spark-sql-streaming-LATEST-with-dependencies.jar",
                "https://repo1.maven.org/maven2/com/redislabs/spark-redis_2.11/2.4.2/spark-redis_2.11-2.4.2-jar-with-dependencies.jar"
            ],
            "python_file_uris": [
                "gs://testing-tmp/proj_BOLSA/"
            ]
        },
    }

    operation = job_client.submit_job_as_operation(
        request={"project_id": project_id, "region": region, "job": job}
    )
    response = operation.result()

    # 解析作业输出路径并获取结果
    matches = re.match("gs://(.*?)/(.*)", response.driver_output_resource_uri)

    output = (
        storage.Client()
        .get_bucket(matches.group(1))
        .blob(f"{matches.group(2)}.000000000")
        .download_as_string()
    )

    print(f"作业执行完成: {output}\r\n")

@staticmethod
def upload_main_file(project, bucket_name, file_name_destination, path_pyspark_file):
    """将PySpark主文件上传到Cloud Storage"""
    print("正在上传PySpark作业文件到Cloud Storage...")
    client = storage.Client(project=project)
    bucket = client.get_bucket(bucket_name)
    blob = bucket.blob(file_name_destination)
    with open(path_pyspark_file, 'rb') as file:
        blob.upload_from_file(file)
    print("PySpark作业文件上传完成。")

@staticmethod
def upload_folder_structures(project, bucket_name, local_path):
    print("正在上传目录结构到Cloud Storage...")
    rel_paths = glob.glob(local_path + '/**', recursive=True)
    client = storage.Client(project=project)
    bucket = client.get_bucket(bucket_name)
    for local_file in rel_paths:
        remote_path = f'{"//".join(local_file.split(os.sep)[settings.NUM_FOLDER:])}'
        if os.path.isfile(local_file):
            blob = bucket.blob(remote_path)
            blob.upload_from_filename(local_file)
    print("目录结构上传完成。")
.....

    def main(self):
        DataProc.upload_main_file(self.project_id, self.bucket_name, self.name_file_transform, self.path_file_transform)
        DataProc.upload_folder_structures(self.project_id, self.bucket_name, self.local_path)
        DataProc.submit_job(self.project_id, self.region, self.cluster_name, self.bucket_name, self.name_file_transform)


if __name__ == '__main__':
    executeTransform = ExecuteTransform(settings.PROJECT_ID, settings.PROJECT_NUMBER, settings.REGION, settings.ZONE,
                                        settings.RESERVATION, settings.TOPIC, settings.SUBSCRIPTION, settings.REGIONAL,
                                        settings.CLUSTER_NAME, settings.BUCKET_NAME, settings.NAME_FILE_TRANSFORM,
                                        settings.PATH_FILE_TRANSFORM, settings.LOCAL_PATH, settings.NUM_FOLDER)
    executeTransform.main()

可能的原因及解决方案

1. Window窗口对齐与Watermark缺失

代码中使用了window(f.col('time'), '30 seconds'),Spark窗口操作默认会对齐到整点时间边界(如00:00:00、00:00:30),且未设置Watermark时,Spark会保留所有历史窗口数据,导致处理逻辑阻塞,触发间隔被拉长。

解决方案:
在process_df_cast中添加Watermark配置,清理过期窗口数据:

df_process_cast = df_process \
    .select('data.*') \
    .withColumn('time', f.to_timestamp(f.col('time'))) \
    .withWatermark('time', '1 minute')  # 根据实际数据延迟调整时长
    # 后续列处理逻辑...

2. Pub/Sub Lite批量拉取配置

Pub/Sub Lite的Spark连接器默认有批量拉取参数(如maxBatchWait、maxBatchSize),若参数设置过大,Spark会等待足够多的消息才触发处理,导致间隔变长。

解决方案:
在read_stream_from_pubsub中添加批量拉取参数,强制控制拉取间隔:

df = (
    spark_session
    .readStream
    .format("pubsublite")
    .option("pubsublite.subscription", f"projects/{project_number}/locations/{region}/subscriptions/{subscription}")
    .option("pubsublite.maxBatchWait", "30s")  # 最多等待30秒拉取批次
    .option("pubsublite.maxBatchSize", "1000")  # 设置较小的批量大小(按需调整)
    .load()
)

3. 单节点集群资源调度延迟

单节点Dataproc集群中,Driver与Executor共享资源,GC停顿、磁盘IO等可能导致触发间隔被拉长。

解决方案:
调整Spark资源配置,启用Streaming指标监控:

spark_session = SparkSession\
    .builder\
    .appName(self.app_name)\
    .master("yarn")\
    .config("spark.executor.memory", "4g")\
    .config("spark.driver.memory", "4g")\
    .config("spark.sql.streaming.metricsEnabled", "true")\
    .getOrCreate()

通过Dataproc监控细粒度的触发时间、处理时间指标,定位延迟环节。

4. Trigger机制的实际生效逻辑

Spark的processingTime触发是上一批次处理完成后等待30秒再触发下一批次,而非严格固定间隔。若某批次处理时间过长,会导致后续批次延迟。

验证方法:
在analysis_batch中添加日志,打印批次执行时间:

def analysis_batch(self, ticker_df_cast, epoch_id):
    import datetime
    print(f"批次 {epoch_id} 执行时间: {datetime.datetime.now()}")

以此确认是触发间隔被拉长,还是批次处理本身耗时过久。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 06:13:10