Untitled
alp
plain_text
10 months ago
6.6 kB
16
Indexable
# kasko_job.py
# Tek dosyada: base_extract -> build_maps -> pivot_write
# spark-submit ile çalıştırın.
import sys
import argparse
from pyspark.sql import SparkSession, functions as F
# =======================
# Ortak ayarlar (DOLDUR)
# =======================
# Iceberg/MinIO katalog ayarlarını Spark tarafında --conf ile de verebilirsiniz.
CATALOG = "iceberg"
SRC_DB = "sompodatabase2"
SRC_TBL = "t001psrcvp"
DATE_COL = "CONFIRM_DATE"
STAGE_DB = "bc_sompo_stage"
BASE_TBL = "kasko_base"
MAP_TBL = "kasko_map"
TARGET_DB = "bc_sompo_test"
TARGET_TBL= "KaskoPivot_Test"
# ---- MSSQL JDBC (DOLDUR) ----
JDBC_URL = "jdbc:sqlserver://<HOST>:1433;databaseName=<DB>;encrypt=true;trustServerCertificate=true"
JDBC_PROPS = {
"user": "<USER>",
"password": "<PASSWORD>",
"driver": "com.microsoft.sqlserver.jdbc.SQLServerDriver",
"fetchsize": "10000"
}
# =======================
# Spark yardımcıları
# =======================
def get_spark(app_name: str) -> SparkSession:
return (
SparkSession.builder
.appName(app_name)
# Gerekirse buraya .config(...) ile MinIO/S3A ve Iceberg conf ekleyin
.getOrCreate()
)
def ensure_table_sql(spark: SparkSession, full_name: str, ddl_cols: str, partition_sql: str = ""):
spark.sql(f"CREATE TABLE IF NOT EXISTS {full_name} ({ddl_cols}) USING ICEBERG {partition_sql}")
# =======================
# 1) BASE EXTRACT
# =======================
def base_extract(process_date_from: str):
spark = get_spark("kasko-base-extract")
src_full = f"{CATALOG}.{SRC_DB}.{SRC_TBL}"
out_full = f"{CATALOG}.{STAGE_DB}.{BASE_TBL}"
df = (spark.table(src_full)
.filter(F.col(DATE_COL) > F.lit(process_date_from)))
ensure_table_sql(
spark,
out_full,
ddl_cols="""
FIRM_CODE INT, COMPANY_CODE INT, PRODUCT_NO STRING,
POLICY_NO STRING, RENEWAL_NO STRING, ENDORS_NO STRING,
QUESTION_CODE STRING, ANSWER STRING, {DATE_COL} DATE, pk_hash BIGINT
""".format(DATE_COL=DATE_COL),
partition_sql=f"PARTITIONED BY (days({DATE_COL}))"
)
(df.writeTo(out_full).append())
print(f"[BASE] {src_full} -> {out_full} (>{process_date_from}) tamam.")
spark.stop()
# =======================
# 2) MAPS (URN SOR -> ÜRÜN -> GRUP)
# =======================
def build_maps():
spark = get_spark("kasko-build-maps")
urnsor_sql = """
SELECT DISTINCT COMPANY_CODE, QUESTION_CODE, PRODUCT_NO
FROM [WINSURE].[fiba].[T001URNSOR] WITH (NOLOCK)
"""
sas_urun_sql = """
SELECT DISTINCT URUN_NO AS U_URUN_NO, URUN_GRUBU_KODU
FROM [WORKDB].[ETLUSR].[SAS_URUN]
"""
sas_urun_grubu_sql = """
SELECT DISTINCT URUN_GRUBU_KODU, URUN_GRUBU_ADI
FROM [WORKDB].[ETLUSR].[SAS_URUN_GRUBU]
"""
urnsor_df = (spark.read.format("jdbc").option("url", JDBC_URL)
.option("query", urnsor_sql).options(**JDBC_PROPS).load())
sas_urun_df = (spark.read.format("jdbc").option("url", JDBC_URL)
.option("query", sas_urun_sql).options(**JDBC_PROPS).load())
sas_urun_grubu_df = (spark.read.format("jdbc").option("url", JDBC_URL)
.option("query", sas_urun_grubu_sql).options(**JDBC_PROPS).load())
map_df = (urnsor_df.alias("u")
.join(sas_urun_df.alias("p"), F.col("u.PRODUCT_NO")==F.col("p.U_URUN_NO"), "inner")
.join(sas_urun_grubu_df.alias("g"), F.col("p.URUN_GRUBU_KODU")==F.col("g.URUN_GRUBU_KODU"), "inner")
.select("u.COMPANY_CODE","u.QUESTION_CODE","u.PRODUCT_NO","g.URUN_GRUBU_ADI")
.dropDuplicates())
out_full = f"{CATALOG}.{STAGE_DB}.{MAP_TBL}"
ensure_table_sql(
spark, out_full,
ddl_cols="COMPANY_CODE INT, QUESTION_CODE STRING, PRODUCT_NO STRING, URUN_GRUBU_ADI STRING"
)
# idempotent yazım
(map_df.writeTo(out_full).overwritePartitions())
print(f"[MAPS] {out_full} güncellendi.")
spark.stop()
# =======================
# 3) PIVOT + TARGET WRITE
# =======================
def sanitize_q_col(colname: str):
q = F.coalesce(F.col(colname).cast("string"), F.lit(""))
q = F.regexp_replace(q, r"[^A-Za-z0-9]+", "_")
q = F.regexp_replace(q, r"_+", "_")
q = F.regexp_replace(q, r"^_|_$", "")
q = F.upper(q)
return F.when(q == "", F.lit("Q_Unkwon")).otherwise(F.concat(F.lit("Q_"), q))
def pivot_write():
spark = get_spark("kasko-pivot-write")
base_full = f"{CATALOG}.{STAGE_DB}.{BASE_TBL}"
map_full = f"{CATALOG}.{STAGE_DB}.{MAP_TBL}"
out_full = f"{CATALOG}.{TARGET_DB}.{TARGET_TBL}"
base = spark.table(base_full)
mapp = spark.table(map_full)
p = (base.select("FIRM_CODE","COMPANY_CODE","PRODUCT_NO",
"POLICY_NO","RENEWAL_NO","ENDORS_NO",
"QUESTION_CODE","ANSWER",DATE_COL,"pk_hash"))
enriched = (p.alias("p")
.join(mapp.alias("m"), ["COMPANY_CODE","QUESTION_CODE","PRODUCT_NO"], "inner")
.filter(F.col("m.URUN_GRUBU_ADI") == F.lit("KASKO"))
.select(F.col("p.pk_hash").alias("ID"),
F.col("p.QUESTION_CODE").alias("QUESTION_CODE"),
F.col("p.ANSWER").alias("ANSWER"))
)
enriched_clean = enriched.withColumn("QCODE", sanitize_q_col("QUESTION_CODE"))
collapsed = (enriched_clean.groupBy("ID","QCODE")
.agg(F.max("ANSWER").alias("ANSWER")))
wide = (collapsed.groupBy("ID").pivot("QCODE").agg(F.first("ANSWER"))
.withColumn("ID", F.col("ID").cast("string")))
(wide.write.format("iceberg").mode("overwrite").saveAsTable(out_full))
print(f"[PIVOT] yazıldı -> {out_full}")
spark.stop()
# =======================
# CLI
# =======================
def parse_args():
p = argparse.ArgumentParser(description="Kasko tek dosya Spark işi")
p.add_argument("--step", choices=["base","maps","pivot","all"], default="all",
help="Hangi adım çalışsın")
p.add_argument("--from-date", default="2025-08-01",
help="Base extract için başlangıç tarihi (YYYY-MM-DD)")
return p.parse_args()
def main():
args = parse_args()
if args.step in ("base", "all"):
base_extract(args.from_date)
if args.step in ("maps", "all"):
build_maps()
if args.step in ("pivot", "all"):
pivot_write()
if __name__ == "__main__":
sys.exit(main())
Editor is loading...
Leave a Comment