Untitled

 avatar
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