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

PyFlink DataStreams API实现24小时手机号-主机访问量统计

必须定义Watermark策略

基于事件时间的窗口聚合必须配置Watermark策略,Kafka流数据几乎都会存在乱序到达的情况,Watermark的作用是告诉Flink:"我认为所有早于某个时间戳的数据都已经到达了",它是Flink处理乱序数据、确定窗口何时可以关闭并计算结果的核心机制。如果不配置,Flink会默认使用处理时间,这和你需求中基于event_time统计过去24小时的逻辑完全不符。

完整实现步骤与代码示例

以下是针对你的场景的完整实现,每一步都做了清晰说明:

  1. 初始化Flink流处理环境
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.connectors import FlinkKafkaConsumer
from pyflink.datastream.window import SlidingEventTimeWindows
from pyflink.common.serialization import SimpleStringSchema
from pyflink.common.typeinfo import Types
from pyflink.common.watermark_strategy import WatermarkStrategy
from pyflink.common.time import Time
from pyflink.datastream.functions import AggregateFunction, MapFunction

# 初始化流处理环境
env = StreamExecutionEnvironment.get_execution_environment()
# 开启checkpoint,故障恢复时避免数据丢失(可选但推荐)
env.enable_checkpointing(30000)  # 每30秒执行一次checkpoint
  1. 配置Kafka数据源
    假设你的Kafka消息是JSON格式(示例:{"phone_number": "13xxxxxxxxx", "host_name": "example.com", "event_time": "2024-05-20 10:30:00"}),后续会解析成Flink可处理的结构:
kafka_consumer = FlinkKafkaConsumer(
    topics="your_topic_name",
    deserialization_schema=SimpleStringSchema(),
    properties={
        "bootstrap.servers": "kafka_broker:9092",
        "group.id": "pyflink_visit_count_group"
    }
)
  1. 定义Watermark策略
    基于event_time字段生成Watermark,同时设置30秒的乱序容忍时间(可根据实际数据流调整):
# 使用Flink内置API快速定义Watermark策略(新手友好)
watermark_strategy = WatermarkStrategy.for_bounded_out_of_orderness(Time.seconds(30)) \
    .with_timestamp_assigner(lambda element, record_timestamp: element[2])
  1. 解析Kafka消息并绑定Watermark
    将JSON字符串解析为包含phone_number、host_name、event_time(毫秒时间戳)的元组:
class JsonParser(MapFunction):
    def map(self, value):
        import json
        from datetime import datetime
        data = json.loads(value)
        # 将字符串格式的event_time转换为毫秒时间戳
        event_ts = int(datetime.strptime(data['event_time'], "%Y-%m-%d %H:%M:%S").timestamp() * 1000)
        return (data['phone_number'], data['host_name'], event_ts)

# 读取Kafka流、解析数据、绑定Watermark策略
stream = env.add_source(kafka_consumer) \
    .map(JsonParser(), output_type=Types.TUPLE([Types.STRING(), Types.STRING(), Types.LONG()])) \
    .assign_timestamps_and_watermarks(watermark_strategy)
  1. 分组并执行窗口聚合
    按phone_number和host_name分组,使用滑动事件时间窗口实现"过去24小时"的实时统计:
# 自定义聚合函数:统计访问次数
class VisitCountAggregate(AggregateFunction):
    def create_accumulator(self):
        return 0
    
    def add(self, value, accumulator):
        return accumulator + 1
    
    def get_result(self, accumulator):
        return accumulator
    
    def merge(self, a, b):
        return a + b

# 按(phone_number, host_name)分组
grouped_stream = stream.key_by(lambda x: (x[0], x[1]))

# 滑动窗口:窗口大小24小时,步长1分钟(步长可根据需求调整)
windowed_stream = grouped_stream.window(SlidingEventTimeWindows.of(Time.hours(24), Time.minutes(1)))

# 应用聚合函数,得到每对(phone_number, host_name)的24小时访问次数
result_stream = windowed_stream.aggregate(VisitCountAggregate(), output_type=Types.LONG())

# 打印结果(也可输出到Kafka、数据库等存储介质)
result_stream.print()

# 启动任务执行
env.execute("Phone-Host 24h Visit Count Job")

关键注意点

  • 如果你的Kafka消息格式不是JSON,需修改JsonParser中的解析逻辑,确保能正确提取三个目标字段。
  • 乱序容忍时间(示例中的30秒)需要根据实际数据流的乱序程度调整,过短会导致部分迟到数据被丢弃,过长会延迟窗口计算结果的输出。
  • 如果只需要按自然天统计(比如每天0点计算前一天的访问次数),可改用TumblingEventTimeWindows.of(Time.hours(24))滚动窗口。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 00:22:22