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

如何用PyFlink优化广告点击聚合输出:仅记录变化及计数器归0

PyFlink广告点击滑动窗口计数优化方案

问题背景

我们有广告点击事件流,需统计每个广告最近5分钟的点击次数,输入输出Schema如下:

输入Schema

ad_id      VARCHAR(10),
clicked_at TIMESTAMP(3)

输出Schema

ad_id      VARCHAR(10),
clicks     INT,
updated_at TIMESTAMP(3)

原滑动窗口方案每分钟输出一次结果,会产生大量重复记录。需要优化为仅当计数实际变化时输出新记录,同时探讨计数归0时如何触发额外输出。


一、仅计数变化时输出的实现

核心思路:先通过滑动窗口计算每个时间点的计数,再对比当前计数与上一次输出的计数,仅保留不一致的记录。

  1. 初始化流环境与输入表
from pyflink.table import EnvironmentSettings, TableEnvironment

env_settings = EnvironmentSettings.in_streaming_mode()
t_env = TableEnvironment.create(env_settings)

# 注册输入Kafka表(可替换为其他数据源)
t_env.execute_sql("""
    CREATE TABLE ad_clicks (
        ad_id VARCHAR(10),
        clicked_at TIMESTAMP(3),
        WATERMARK FOR clicked_at AS clicked_at - INTERVAL '5' SECOND
    ) WITH (
        'connector' = 'kafka',
        'topic' = 'ad_click_topic',
        'properties.bootstrap.servers' = 'localhost:9092',
        'format' = 'json'
    )
""")
  1. 滑动窗口聚合计算5分钟窗口内的点击数(每分钟更新一次)
windowed_click_stats = t_env.sql_query("""
    SELECT
        ad_id,
        COUNT(*) AS clicks,
        TUMBLE_END(clicked_at, INTERVAL '1' MINUTE) AS updated_at
    FROM ad_clicks
    GROUP BY
        ad_id,
        TUMBLE(clicked_at, INTERVAL '5' MINUTE, INTERVAL '1' MINUTE)
""")
  1. 过滤仅计数变化的记录
    用LAG函数获取每个广告上一次输出的计数,仅当当前计数与上次不同时输出:
final_result = t_env.sql_query("""
    SELECT
        ad_id,
        clicks,
        updated_at
    FROM (
        SELECT
            ad_id,
            clicks,
            updated_at,
            LAG(clicks) OVER (PARTITION BY ad_id ORDER BY updated_at) AS prev_clicks
        FROM %s
    )
    WHERE prev_clicks IS NULL OR clicks != prev_clicks
""" % windowed_click_stats)
  1. 注册输出表并执行
t_env.execute_sql("""
    CREATE TABLE ad_click_stats_output (
        ad_id VARCHAR(10),
        clicks INT,
        updated_at TIMESTAMP(3)
    ) WITH (
        'connector' = 'kafka',
        'topic' = 'ad_click_stats_topic',
        'properties.bootstrap.servers' = 'localhost:9092',
        'format' = 'json'
    )
""")

final_result.execute_insert("ad_click_stats_output").wait()

二、计数归0时输出额外记录

滑动窗口默认仅在有事件触发时计算计数,当窗口内所有事件过期后,计数归0但不会自动触发输出。要实现此需求,需结合状态编程+定时器,通过自定义ProcessTableFunction处理:

核心逻辑

  • 维护每个广告的状态:窗口内有效点击事件列表、当前计数、上一次输出的计数
  • 为每个点击事件设置定时器,在事件过期时间(clicked_at + 5分钟)触发
  • 定时器触发时,清理过期事件并更新计数,若计数从非0变为0,则输出归0记录

实现代码

from pyflink.table import DataTypes
from pyflink.table.udf import ProcessTableFunction
from pyflink.datastream import TimerService
from collections import defaultdict
import datetime

class AdClickCountProcess(ProcessTableFunction):
    def open(self, context):
        # 每个广告的有效事件时间列表
        self.ad_events = defaultdict(list)
        # 每个广告的当前计数与上一次输出计数
        self.ad_count_state = defaultdict(lambda: {"current": 0, "last_output": 0})
        self.timer_service: TimerService = context.timer_service()

    def process_element(self, row):
        ad_id = row.ad_id
        clicked_ts = row.clicked_at.timestamp() * 1000  # 转为毫秒级时间戳
        # 更新事件列表与当前计数
        self.ad_events[ad_id].append(clicked_ts)
        self.ad_count_state[ad_id]["current"] += 1
        # 设置事件过期定时器(5分钟后)
        self.timer_service.register_event_time_timer(clicked_ts + 5 * 60 * 1000)
        # 仅当当前计数与上次输出不同时,输出记录
        if self.ad_count_state[ad_id]["current"] != self.ad_count_state[ad_id]["last_output"]:
            yield (
                ad_id,
                self.ad_count_state[ad_id]["current"],
                row.clicked_at
            )
            self.ad_count_state[ad_id]["last_output"] = self.ad_count_state[ad_id]["current"]

    def on_timer(self, timestamp, ctx):
        # 遍历所有广告,清理过期事件并检查计数变化
        for ad_id in list(self.ad_events.keys()):
            # 过滤出最近5分钟内的有效事件
            valid_ts = [ts for ts in self.ad_events[ad_id] if ts > timestamp - 5 * 60 * 1000]
            new_count = len(valid_ts)
            last_output = self.ad_count_state[ad_id]["last_output"]
            # 若计数从非0变为0,输出归0记录
            if last_output != 0 and new_count == 0:
                yield (
                    ad_id,
                    0,
                    datetime.datetime.fromtimestamp(timestamp / 1000)
                )
                self.ad_count_state[ad_id]["last_output"] = 0
            # 更新状态
            self.ad_events[ad_id] = valid_ts
            self.ad_count_state[ad_id]["current"] = new_count

# 注册自定义处理函数
t_env.register_function("ad_click_count_process", AdClickCountProcess())

# 使用自定义函数处理流
result_with_zero = t_env.sql_query("""
    SELECT
        ad_id,
        clicks,
        updated_at
    FROM ad_clicks,
    LATERAL TABLE(ad_click_count_process(ad_id, clicked_at)) AS T(ad_id, clicks, updated_at)
""")

# 输出到目标表(可复用之前的输出表)
result_with_zero.execute_insert("ad_click_stats_output").wait()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 00:25:26