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

