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))

更多推荐