Note:- I am using pyflink lient api and not submit...
# random
m
Note:- I am using pyflink lient api and not submitting the job to flink cluster also I have added the flink_kafka_connector jar as per below
Copy code
# add kafka connector dependency
kafka_jar = os.path.join(os.path.abspath(os.path.dirname(__file__)),
                         'external_libs\\flink-sql-connector-kafka_2.11-1.13.0.jar')

print("FLAG , {}".format(kafka_jar))
tbl_env.get_config() \
    .get_configuration() \
    .set_string("pipeline.jars", "file://{}".format(kafka_jar))
Flink version is 1.13.0