flink执行sql文件代码
·
from pyflink.table import EnvironmentSettings, TableEnvironment
import os
sql_file = "/data/dba4/flink-1.18.0/task/cdc.sql"
# # 1. create a TableEnvironment
env_settings = EnvironmentSettings.in_streaming_mode()
table_env = TableEnvironment.create(env_settings)
def read_sqls(sql_path):
with open(sql_path, 'r') as f:
sql = f.read()
return sql.split(';')
def execute_sql(sql):
if not sql or len(sql) < 10:
return
print(sql)
table_env.execute_sql(sql)
for sql in read_sqls(sql_file):
try:
print("execute sql: {}".format(sql))
execute_sql(sql)
except Exception as e:
print("error: {}".format(e))
更多推荐


所有评论(0)