-- Databricks notebook source -- ============================================================================= -- Purpose : Daily ingestion of the Pharbers provincial sales fact file -- (Pharbers_PROV_Fact*.csv) from the ADLS user-upload blob into -- the DWD province-level sales fact. -- Source : ADLS blob CSV files (Pharbers_PROV_Fact*.csv under -- <...>/ODS/GND/UserUpload//, environment-specific), -- tmp.tmp_chpa_raw_data (staging) -- Target : dwd.dwd_gnd_pharbers_prov_fact -- Grain : One row per province x month x drug as provided by the source -- file; grain to be confirmed against the CSV contract. -- Write mode : Full refresh (INSERT OVERWRITE). -- Replaces : Legacy CHPA/01_FB_BLOB_TO_DWD.sql (same DWD target). -- Consumers : 02 tmp_ims_tf_fact_sales -> DM sales. -- Notes : Two engineering fixes vs legacy (per migration doc): -- 1) Fail fast when the day's upload path or matching files are -- missing, instead of silently republishing stale staging. -- 2) withColumnRenamed results are now assigned back to the frame -- (legacy discarded them); renames are applied only to columns -- that exist, preserving legacy tolerance of missing columns. -- Everything else is preserved: paths, file filter regex, -- CSV options, staging flow, NULL lineage columns, and the UTC+8 -- insert timestamp. -- ============================================================================= -- COMMAND ---------- -- MAGIC %run ../../../Common/config -- COMMAND ---------- -- MAGIC %md -- MAGIC ### 从 blob 读取 csv 文件作为 CHPA 法伯省级事实表 -- COMMAND ---------- -- MAGIC %python -- MAGIC from datetime import datetime, timedelta -- MAGIC import pandas as pd -- COMMAND ---------- -- MAGIC %python -- MAGIC if ENVIRONMENT == PRD_ENVIRONMENT_VALUE: -- MAGIC factsales_file_path_template = "abfss://master@azcdatalakeprd.dfs.core.chinacloudapi.cn/ODS/GND/UserUpload/" -- MAGIC elif ENVIRONMENT == TEST_ENVIRONMENT_VALUE: -- MAGIC factsales_file_path_template = "abfss://master@retaildlstoragetest.dfs.core.chinacloudapi.cn/ODS/GND/UserUpload/" -- COMMAND ---------- -- MAGIC %python -- MAGIC # 路径是否存在 -- MAGIC def path_exists(path): -- MAGIC try: -- MAGIC dbutils.fs.ls(path) -- MAGIC return True -- MAGIC except Exception as e: -- MAGIC if "java.io.FileNotFoundException" in str(e): -- MAGIC return False -- MAGIC else: -- MAGIC print(f"检查路径 {path} 时出错: {e}") -- MAGIC raise -- COMMAND ---------- -- MAGIC %python -- MAGIC # 列出 blob 上的文件列表 -- MAGIC def list_file_name(path): -- MAGIC first_path_list = [i.path for i in dbutils.fs.ls(path)] -- MAGIC second_path_list = [dbutils.fs.ls(i)[0] for i in first_path_list] -- MAGIC return second_path_list -- COMMAND ---------- -- MAGIC %python -- MAGIC # 从 blob 下载文件到 local -- MAGIC def download_file(file_path, local_path): -- MAGIC dbutils.fs.cp(file_path, local_path) -- MAGIC print(f"已下载 {file_path} 到 {local_path}") -- MAGIC return local_path -- COMMAND ---------- -- MAGIC %python -- MAGIC # 计算时间得到当天的路径 -- MAGIC current_date = datetime.utcnow() + timedelta(hours=8) -- MAGIC date_path = current_date.strftime("%Y/%m/%d/") -- MAGIC base_path = factsales_file_path_template + date_path -- COMMAND ---------- -- MAGIC %md -- MAGIC ### 获取路径下的文件名称,并挑出符合条件的文件路径 -- MAGIC - 无文件时直接失败,避免旧 staging 数据被再次覆盖到 DWD -- COMMAND ---------- -- MAGIC %python -- MAGIC if not path_exists(base_path): -- MAGIC raise FileNotFoundError(f"上传路径不存在,任务失败,拒绝沿用旧 staging 数据: {base_path}") -- MAGIC -- MAGIC all_file_list = list_file_name(base_path) -- MAGIC -- MAGIC # 生成 df 来筛选内容 -- MAGIC files_df = pd.DataFrame([{ -- MAGIC 'path': f.path, -- MAGIC 'modificationtime': f.modificationTime, -- MAGIC 'name': f.name -- MAGIC } for f in all_file_list]) -- MAGIC -- MAGIC # 同名文件保留修改时间最新的一份 -- MAGIC files_df = files_df.sort_values('modificationtime', ascending=False).drop_duplicates('name').sort_index() -- MAGIC files_df = files_df[files_df['name'].str.match(r'^Pharbers_PROV_Fact.*\.csv$')] -- MAGIC files_df = files_df.reset_index(drop=True) -- MAGIC -- MAGIC if files_df.empty: -- MAGIC raise RuntimeError("未找到符合条件的数据文件 (Pharbers_PROV_Fact*.csv),任务失败,拒绝沿用旧 staging 数据") -- MAGIC -- MAGIC print(f"找到 {len(files_df)} 个符合条件的数据文件") -- COMMAND ---------- -- MAGIC %python -- MAGIC import os -- MAGIC -- MAGIC # 下载数据到 local,读取并清洗,逐文件收集 -- MAGIC df_all = [] -- MAGIC for file in files_df['path'].tolist(): -- MAGIC local_path = download_file(file, f"/Volumes/{NGBI_CATALOG}/tmp/volume_tmp/tmp/{os.path.basename(file)}") -- MAGIC file_df = (spark.read -- MAGIC .option("header", "true") -- MAGIC .option("quote", '"') -- MAGIC .option("escape", '"') -- MAGIC .option("multiLine", "true") -- MAGIC .option("mode", "PERMISSIVE") -- MAGIC .csv(local_path)) -- MAGIC # 与旧脚本一致:丢弃 TA / Market 两列 -- MAGIC file_df = file_df.drop("TA", "Market") -- MAGIC -- MAGIC # 修复历史 bug:旧脚本调用 withColumnRenamed 后未把结果赋回原变量, -- MAGIC # 重命名实际从未生效。这里正确赋值;且仅对确实存在的列重命名, -- MAGIC # 保持旧任务对缺失列名的容忍,不因列不存在而失败。 -- MAGIC rename_map = { -- MAGIC 'IMS.药品ID': 'IMS_DRUG_ID', -- MAGIC '是否法伯编码': 'IS_HOSP_CODE', -- MAGIC '规格': 'SPEC', -- MAGIC '转换比': 'CONVERSION_RATIO', -- MAGIC '剂型': 'DOSAGE_FORM', -- MAGIC '价格': 'PRICE', -- MAGIC } -- MAGIC for old_name, new_name in rename_map.items(): -- MAGIC if old_name in file_df.columns: -- MAGIC file_df = file_df.withColumnRenamed(old_name, new_name) -- MAGIC -- MAGIC print(f"已读取 {local_path}") -- MAGIC df_all.append(file_df) -- MAGIC -- MAGIC if not df_all: -- MAGIC raise RuntimeError("没有读取到任何数据文件,任务失败") -- COMMAND ---------- -- MAGIC %python -- MAGIC # 先清空 staging,避免旧数据残留,随后逐文件追加 -- MAGIC spark.sql("TRUNCATE TABLE tmp.tmp_chpa_raw_data") -- MAGIC for num, file_df in enumerate(df_all, start=1): -- MAGIC file_df.createOrReplaceTempView("fact_sales") -- MAGIC spark.sql("INSERT INTO tmp.tmp_chpa_raw_data SELECT * FROM fact_sales") -- MAGIC print(f"第{num}个") -- COMMAND ---------- -- 全量覆盖 INSERT OVERWRITE TABLE dwd.dwd_gnd_pharbers_prov_fact ( year, ym, province_c, ims_drug_id, is_hosp_code, prod_corp, prod_cod, pack_cod, phcd, prod_des, cmps_des, corp_des, mnfl_cod, prod_des_c, cmps_c, corp_des_c, pack_des, spec, conversion_ratio, dosage_form, atc4_cod, app1_cod, app1_des, app1_des_c, app2_cod, app2_des, app2_des_c, app3_cod, app3_des, app3_des_c, vbp_batch, vbp, value, totalunit, countingunit, price, manu_des, manu_des_c, source_file_path, source_file_name, etl_insert_dt ) SELECT year, ym, province_c, ims_drug_id, is_hosp_code, prod_corp, prod_cod, pack_cod, phcd, prod_des, cmps_des, corp_des, mnfl_cod, prod_des_c, cmps_c, corp_des_c, pack_des, spec, conversion_ratio, dosage_form, atc4_cod, app1_cod, app1_des, app1_des_c, app2_cod, app2_des, app2_des_c, app3_cod, app3_des, app3_des_c, vbp_batch, vbp, value, totalunit, countingunit, price, manu_des, manu_des_c, NULL AS source_file_path, NULL AS source_file_name, FROM_UTC_TIMESTAMP(CURRENT_TIMESTAMP(), 'UTC+8') AS etl_insert_dt FROM tmp.tmp_chpa_raw_data ;