在K8s上的Flink集群提交Python作业时卡住
K8s环境下Flink Python作业提交失败排查
问题现象
- 示例Python作业可正常执行:
./bin/flink run --python examples/python/table/batch/word_count.py - 提交自定义Python作业时控制台卡住,Job Manager中无该作业记录,控制台输出日志如下:
./bin/flink run --python bin/flink.py ERROR StatusLogger Reconfiguration failed: No configuration found for '5ce65a89' at 'null' in 'null' WARNING: An illegal reflective access operation has occurred WARNING: Illegal reflective access by org.apache.flink.api.java.ClosureCleaner (file:/opt/flink/lib/apache-flink-examples-0.0.2-RC3.jar) to field java.util.Properties.serialVersionUID WARNING: Please consider reporting this to the maintainers of org.apache.flink.api.java.ClosureCleaner WARNING: Use --illegal-access=warn to enable warnings of further illegal reflective access operations WARNING: All illegal access operations will be denied in a future release - 功能等价的Java作业可正常运行。
自定义作业代码(flink.py)
import logging import sys from pyflink.table import (EnvironmentSettings, TableEnvironment, DataTypes, StreamTableEnvironment) from pyflink.datastream import StreamExecutionEnvironment def main(): env = StreamExecutionEnvironment.get_execution_environment() env.add_jars("file:///opt/lib/flink-connector-kafka_2.11-1.13.6.jar;file:///flink-json-1.13.6.jar,file:///flink-sql-connector-kafka_2.11-1.13.6.jar") t_env = StreamTableEnvironment.create(env) # source t_env.execute_sql(""" CREATE TABLE trades ( instrumentId INT, exchange STRING, category1 STRING, category2 STRING, traderId INT, price DOUBLE, quantity DOUBLE ) WITH ( 'connector' = 'kafka', 'topic' = 'jse-trades', 'properties.bootstrap.servers' = '<kafka_url>', 'properties.group.id' = 'flink-reader', 'scan.startup.mode' = 'earliest-offset', 'format' = 'json' ) """) # sink t_env.execute_sql(""" CREATE TABLE trades_aggr ( instrumentId INT, exchange STRING, category1 STRING, category2 STRING, traderId INT, trade_count BIGINT, trade_turnover DOUBLE, trade_volume DOUBLE, trade_volume_avg DOUBLE, PRIMARY KEY (instrumentId) NOT ENFORCED ) WITH ( 'connector' = 'upsert-kafka', 'topic' = 'jse-trades-results', 'properties.bootstrap.servers' = '<kafka_url>', 'properties.group.id' = 'testGroup', 'sink.buffer-flush.max-rows' = '5000', 'sink.buffer-flush.interval' = '1000', 'sink.parallelism' = '2', 'key.format' = 'json', 'value.format' = 'json' ) """) # aggregation t_env.execute_sql(""" INSERT INTO trades_aggr SELECT instrumentId, exchange, category1, category2, traderId, count(*), sum(price*quantity), sum(quantity), avg(quantity) FROM trades GROUP BY instrumentId, exchange, category1, category2, traderId """) t_env.execute("kafka_trades") if __name__ == "__main__": logging.basicConfig(stream=sys.stdout, level=logging.INFO, format="%(message)s") main()
需排查自定义Python作业提交失败的原因,确认是否遗漏必要配置或步骤。
内容的提问来源于stack exchange,提问作者Ardelia Lortz
相关产品推荐
相关产品推荐

