from clickzetta_dbutils import get_active_lakehouse_engine
from sqlalchemy import text
engine = get_active_lakehouse_engine(schema="your_schema")
with engine.connect() as conn:
conn.execute(text("SELECT 1"))
场景:已有 shell 脚本接入调度
典型场景:团队有一批用
awk
awk
/
sed
sed
处理日志或 CSV 的 shell 脚本,想直接接入 Studio 调度体系,处理完后把结果写入 Lakehouse。
python3 -c "
import json
posts = json.load(open('/tmp/posts.json'))
for p in posts:
print(f\"{p['id']},{p['userId']},{p['title'][:30].replace(',','')}\")
" | awk -F, '$2 <= 3 {print}' > /tmp/posts_filtered.csv
echo "过滤后行数:$(wc -l < /tmp/posts_filtered.csv)"
用 python3 把结果写入 Lakehouse:
from clickzetta_dbutils import get_active_lakehouse_engine
from sqlalchemy import text
biz_date = '$BIZ_DATE'
engine = get_active_lakehouse_engine(schema="doc_connector_demo")
with engine.connect() as conn:
conn.execute(text("CREATE SCHEMA IF NOT EXISTS doc_connector_demo"))
conn.execute(text("""
CREATE TABLE IF NOT EXISTS doc_connector_demo.doc_shell_posts (
post_id INT,
user_id INT,
title STRING,
load_date STRING
)
"""))
conn.execute(text(f"DELETE FROM doc_connector_demo.doc_shell_posts WHERE load_date = '{biz_date}'"))
rows = 0
with open('/tmp/posts_filtered.csv') as f:
for line in f:
parts = line.strip().split(',', 2)
if len(parts) == 3:
post_id, user_id, title = parts
title = title.replace("'", "''")
conn.execute(text(
f"INSERT INTO doc_connector_demo.doc_shell_posts VALUES "
f"({post_id}, {user_id}, '{title}', '{biz_date}')"
))
rows += 1
print(f"写入 {rows} 行,load_date={biz_date}")
with engine.connect() as conn:
result = conn.execute(text(
f"SELECT COUNT(*) as cnt, COUNT(DISTINCT user_id) as users "
f"FROM doc_connector_demo.doc_shell_posts WHERE load_date = '{biz_date}'"
))
row = result.fetchone()
print(f"验证:{row[0]} 条记录,{row[1]} 个用户")