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

Apache Flink无法打印新Kafka流数据 考勤场景求助

问题:Flink处理Kafka门禁日志时新用户无输出,老用户正常

我正在做办公室用户门禁日志的实时考勤数据处理:用户刷卡后日志事件进入Kafka,消费后要判断用户是否首次访问,首次则标记考勤时间,非首次则增加访问次数。

作为Flink新手,我写了如下Python代码,但遇到问题:发送新用户(如'Bhuvi20')的日志时Flink无任何输出;发送已处理过的用户(如'Bhuvi19')数据时能正常打印,新数据发送后程序好像停滞了。

import json
import os

from pyflink.common import SimpleStringSchema
from pyflink.datastream import StreamExecutionEnvironment, RuntimeExecutionMode
from pyflink.datastream.connectors import FlinkKafkaConsumer
from pyflink.common.typeinfo import Types
from pyflink.table import StreamTableEnvironment, EnvironmentSettings


def my_map(obj):
    json_obj = json.loads(json.loads(obj))
    return json.dumps(json_obj["name"])


def kafkaread():
    env = StreamExecutionEnvironment.get_execution_environment()
    env.set_runtime_mode(RuntimeExecutionMode.STREAMING)

    jars_path = "/home/vishal/flink/flink-pipeline/flink-env3.10/flink_test_1/"
    jar_files = [
        "file:///" + jars_path + "flink-sql-connector-kafka-3.1.0-1.18.jar",
    ]
    jar_files_str = ";".join(jar_files)

    env.add_jars(jar_files_str)

    t_env = StreamTableEnvironment.create(stream_execution_environment=env)
    t_env.get_config().set('table.exec.source.idle-timeout', '1000 ms')

   
    # with below code I am taking the latest data which is coming in kafka
    create_table_sql = """
        CREATE TABLE access_logs_source (
                        `name` STRING,
                        `door` STRING,
                        `access_time` TIMESTAMP
                      ) WITH (
                        'connector' = 'kafka',
                        'topic' = 'access_logs_topic_5',
                        'properties.bootstrap.servers' = 'localhost:9092',
                        'properties.group.id' = 'flink-group-access_logs_5',
                        'scan.startup.mode' = 'latest-offset',
                        'format' = 'json'
                      );
      """
    # It will fetch all the data which is present in kafka along with latest data so that I can check if latest data is present here or not
    historical_table_sql = """
        CREATE TABLE access_logs_historical (
                        `name` STRING,
                        `door` STRING,
                        `access_time` TIMESTAMP
                      ) WITH (
                        'connector' = 'kafka',
                        'topic' = 'access_logs_topic_5',
                        'properties.bootstrap.servers' = 'localhost:9092',
                        'properties.group.id' = 'flink-group-access_logs_5',
                        'scan.startup.mode' = 'earliest-offset',
                        'format' = 'json'
                      );
      """

    # Execute the SQL command to create the table
    t_env.execute_sql(create_table_sql)
    
    t_env.execute_sql(historical_table_sql)

    # Define the SQL command to select data from the created table
    select_sql = "SELECT name FROM access_logs_source"

    # Execute the SQL command to select data from the table and print the result
    result_table = t_env.execute_sql(select_sql)
    # result_table.print()
    result_rows = result_table.collect()
    # new log data of the user
    name_to_filter = ''
    for row in result_rows:
        # print(row[0])

        name_to_filter = row[0]
        # break
        print('-------', name_to_filter)
            # result_rows.close()
            
        # checking if new log data is present in the kafka topic or not
        check_query = """
          SELECT * from access_logs_historical where name = '{0}'
        """
        check_query = check_query.format(name_to_filter)
        print('++++++++++ ', check_query)
        # t_env.execute_sql(check_query)
        table = t_env.sql_query(check_query)
        table_result = table.execute()
        table_result.print()
    

    # Execute the Flink job
    # t_env.execute("Create and Select Table Job")

if __name__ == '__main__':
    kafkaread()

问题分析与修复

核心问题点

  1. 错误的历史数据查询方式:用两个Kafka表(latest/earliest偏移量)来判断用户是否存在,完全违背Flink流处理范式。Kafka是无限流,earliest-offset的表会一直从头消费,永远不会结束,导致新用户查询时程序阻塞。
  2. 阻塞式的collect()调用:result_rows = result_table.collect()在流模式下是阻塞操作,循环内嵌套执行新SQL查询会彻底卡死线程,无法处理后续新数据。
  3. 重复消费冲突:两个表用同一个group.id,会导致Kafka消费位移混乱,无法正确跟踪消费进度。
  4. 语法错误:代码中使用//作为注释,Python中应使用#,会引发语法问题。

修复后的代码实现

正确的做法是利用Flink的状态管理自动维护用户访问历史,用SQL聚合直接计算首次访问时间和访问次数:

import os
from pyflink.table import StreamTableEnvironment, EnvironmentSettings

def process_access_logs():
    # 初始化流处理环境
    env_settings = EnvironmentSettings.in_streaming_mode()
    t_env = StreamTableEnvironment.create(environment_settings=env_settings)

    # 设置Kafka连接器jar包路径
    jars_path = "/home/vishal/flink/flink-pipeline/flink-env3.10/flink_test_1/"
    jar_files = [
        "file:///" + jars_path + "flink-sql-connector-kafka-3.1.0-1.18.jar",
    ]
    t_env.get_config().set("pipeline.jars", ";".join(jar_files))

    # 创建Kafka源表(添加水印处理时间乱序)
    create_source_sql = """
        CREATE TABLE access_logs_source (
            `name` STRING,
            `door` STRING,
            `access_time` TIMESTAMP(3),
            WATERMARK FOR access_time AS access_time - INTERVAL '5' SECOND
        ) WITH (
            'connector' = 'kafka',
            'topic' = 'access_logs_topic_5',
            'properties.bootstrap.servers' = 'localhost:9092',
            'properties.group.id' = 'flink-group-access-logs',
            'scan.startup.mode' = 'latest-offset',
            'format' = 'json'
        );
    """
    t_env.execute_sql(create_source_sql)

    # 聚合计算每个用户的访问信息
    process_sql = """
        SELECT 
            name,
            MIN(access_time) AS first_access_time,
            COUNT(*) AS access_count,
            CASE WHEN COUNT(*) = 1 THEN '首次访问' ELSE '重复访问' END AS access_type
        FROM access_logs_source
        GROUP BY name;
    """

    # 执行查询并实时输出结果
    result_table = t_env.sql_query(process_sql)
    result_table.execute().print()

if __name__ == '__main__':
    process_access_logs()

修复说明

  • 状态化聚合:GROUP BY name让Flink自动为每个用户维护状态,无需手动扫描历史数据,首次访问时COUNT(*)为1,后续递增。
  • 去掉阻塞操作:直接由Flink流引擎处理数据,避免collect()和嵌套查询导致的线程阻塞。
  • 水印机制:添加水印处理门禁日志可能出现的时间乱序问题,确保计算结果准确。
  • 单一消费者:只用一个Kafka消费者表,避免位移冲突,简化逻辑。

原代码中新用户无输出的原因:当查询新用户时,access_logs_historical表从最早偏移量开始消费,而Kafka是无限流,该表永远不会消费完成,导致table_result.print()一直阻塞,程序停滞。老用户数据已存在于主题中,能快速返回结果,所以正常打印。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 00:49:51