Files
airflow-coolify/scripts/bigquery_aggraget_fact_selected_layer.py
T

1205 lines
57 KiB
Python

"""
BIGQUERY ANALYSIS LAYER - INDICATOR NORM AGGREGATION
PERUBAHAN ARSITEKTUR:
- ASEAN aggregate DIGABUNG ke dalam tabel yang sama (country_id=0, country_name="ASEAN")
sehingga Looker Studio dapat memfilter: per negara, atau ASEAN saja.
- agg_narrative_indicator: granularity tetap per indicator_id (all years, all countries),
ASEAN summary ditambahkan sebagai kolom terpisah (asean_avg_value_first/last).
Output 2 tabel:
1. agg_indicator_norm -> per baris (year x country x indicator), termasuk ASEAN rows
2. agg_narrative_indicator -> per indicator_id, ada kolom asean_avg_* tambahan
BUGFIX (diteruskan dari versi sebelumnya):
- INDICATOR_NAME_ID_MAP: semua key lowercase agar match dengan .lower().strip() lookup.
"""
import pandas as pd
import numpy as np
from datetime import datetime
import logging
import json
from scripts.bigquery_config import get_bigquery_client
from scripts.bigquery_helpers import (
log_update,
load_to_bigquery,
read_from_bigquery,
setup_logging,
save_etl_metadata,
)
from google.cloud import bigquery
# =============================================================================
# KONSTANTA
# =============================================================================
ASEAN_COUNTRY_ID = 0
ASEAN_COUNTRY_NAME = "ASEAN"
ASEAN_COUNTRY_NAME_ID = "ASEAN"
# =============================================================================
# MAPPING BAHASA INDONESIA
# CHANGED: Other / Lainnya
# =============================================================================
COUNTRY_NAME_ID_MAP: dict = {
"Brunei Darussalam" : "Brunei Darussalam",
"Cambodia" : "Kamboja",
"Indonesia" : "Indonesia",
"Lao People's Democratic Republic" : "Laos",
"Lao PDR" : "Laos",
"Malaysia" : "Malaysia",
"Myanmar" : "Myanmar",
"Philippines" : "Filipina",
"Singapore" : "Singapura",
"Thailand" : "Thailand",
"Timor-Leste" : "Timor-Leste",
"Viet Nam" : "Vietnam",
"Vietnam" : "Vietnam",
"ASEAN" : "ASEAN",
}
PILLAR_NAME_ID_MAP: dict = {
# Mapping nama pilar (Inggris dengan prefix Food) -> Bahasa Indonesia
"Food Availability" : "Ketersediaan Pangan",
"Food Access" : "Akses Pangan",
"Food Utilization" : "Pemanfaatan Pangan",
"Food Stability" : "Stabilitas Pangan",
"Food Other" : "Indikator Tambahan",
# Variasi tanpa prefix Food
"Availability" : "Ketersediaan Pangan",
"Access" : "Akses Pangan",
"Utilization" : "Pemanfaatan Pangan",
"Stability" : "Stabilitas Pangan",
"Other" : "Indikator Tambahan",
# lowercase
"food availability" : "Ketersediaan Pangan",
"food access" : "Akses Pangan",
"food utilization" : "Pemanfaatan Pangan",
"food stability" : "Stabilitas Pangan",
"food other" : "Indikator Tambahan",
"availability" : "Ketersediaan Pangan",
"access" : "Akses Pangan",
"utilization" : "Pemanfaatan Pangan",
"stability" : "Stabilitas Pangan",
"other" : "Indikator Tambahan",
}
# BUGFIX: semua key lowercase
INDICATOR_NAME_ID_MAP: dict = {
"dietary energy supply used in the estimation of the prevalence of undernourishment (kcal/cap/day)":
"Pasokan energi makanan yang digunakan dalam estimasi prevalensi kekurangan gizi (kkal/kapita/hari)",
"dietary energy supply used in the estimation of the prevalence of undernourishment (kcal/cap/day) (3-year average)":
"Pasokan energi makanan yang digunakan dalam estimasi prevalensi kekurangan gizi (kkal/kapita/hari) (rata-rata 3 tahun)",
"percentage of population using at least basic drinking water services (percent)":
"Persentase penduduk yang menggunakan layanan air minum dasar (persen)",
"percentage of population using at least basic sanitation services (percent)":
"Persentase penduduk yang menggunakan layanan sanitasi dasar (persen)",
"percentage of population using safely managed drinking water services (percent)":
"Persentase penduduk yang menggunakan layanan air minum yang dikelola dengan aman (persen)",
"percentage of population using safely managed sanitation services (percent)":
"Persentase penduduk yang menggunakan layanan sanitasi yang dikelola dengan aman (persen)",
"rail lines density (total route in km per 100 square km of land area)":
"Kepadatan jalur kereta api (total rute dalam km per 100 km² lahan)",
"average dietary energy requirement (kcal/cap/day)":
"Rata-rata kebutuhan energi makanan (kkal/kapita/hari)",
"average dietary energy supply adequacy (percent) (3-year average)":
"Kecukupan rata-rata pasokan energi makanan (persen) (rata-rata 3 tahun)",
"average fat supply (g/cap/day) (3-year average)":
"Rata-rata pasokan lemak (g/kapita/hari) (rata-rata 3 tahun)",
"average protein supply (g/cap/day) (3-year average)":
"Rata-rata pasokan protein (g/kapita/hari) (rata-rata 3 tahun)",
"average supply of protein of animal origin (g/cap/day) (3-year average)":
"Rata-rata pasokan protein hewani (g/kapita/hari) (rata-rata 3 tahun)",
"percent of arable land equipped for irrigation (percent) (3-year average)":
"Persentase lahan pertanian yang dilengkapi irigasi (persen) (rata-rata 3 tahun)",
"cereal import dependency ratio (percent) (3-year average)":
"Rasio ketergantungan impor sereal (persen) (rata-rata 3 tahun)",
"share of dietary energy supply derived from cereals, roots and tubers (percent) (3-year average)":
"Proporsi pasokan energi makanan dari serealia, akar, dan umbi-umbian (persen) (rata-rata 3 tahun)",
"per capita food supply variability (kcal/cap/day)":
"Variabilitas pasokan pangan per kapita (kkal/kapita/hari)",
"value of food imports in total merchandise exports (percent) (3-year average)":
"Nilai impor pangan terhadap total ekspor barang (persen) (rata-rata 3 tahun)",
"gross domestic product per capita, ppp, (constant 2021 international $)":
"Produk domestik bruto per kapita, PPP (internasional konstan 2021 US$)",
"political stability and absence of violence/terrorism (index)":
"Stabilitas politik dan tidak adanya kekerasan/terorisme (indeks)",
"prevalence of undernourishment (percent) (3-year average)":
"Prevalensi kekurangan gizi (persen) (rata-rata 3 tahun)",
"number of people undernourished (million) (3-year average)":
"Jumlah penduduk kekurangan gizi (juta jiwa) (rata-rata 3 tahun)",
"minimum dietary energy requirement (kcal/cap/day)":
"Kebutuhan energi makanan minimum (kkal/kapita/hari)",
"prevalence of exclusive breastfeeding among infants 0-5 months of age (percent)":
"Prevalensi pemberian ASI eksklusif pada bayi usia 0-5 bulan (persen)",
"number of children under 5 years affected by wasting (million)":
"Jumlah anak di bawah 5 tahun yang mengalami wasting (juta jiwa)",
"number of moderately or severely food insecure female adults (million) (3-year average)":
"Jumlah perempuan dewasa yang mengalami kerawanan pangan sedang atau berat (juta jiwa) (rata-rata 3 tahun)",
"number of moderately or severely food insecure male adults (million) (3-year average)":
"Jumlah laki-laki dewasa yang mengalami kerawanan pangan sedang atau berat (juta jiwa) (rata-rata 3 tahun)",
"number of moderately or severely food insecure people (million) (3-year average)":
"Jumlah penduduk yang mengalami kerawanan pangan sedang atau berat (juta jiwa) (rata-rata 3 tahun)",
"number of severely food insecure female adults (million) (3-year average)":
"Jumlah perempuan dewasa yang mengalami kerawanan pangan berat (juta jiwa) (rata-rata 3 tahun)",
"number of severely food insecure male adults (million) (3-year average)":
"Jumlah laki-laki dewasa yang mengalami kerawanan pangan berat (juta jiwa) (rata-rata 3 tahun)",
"number of severely food insecure people (million) (3-year average)":
"Jumlah penduduk yang mengalami kerawanan pangan berat (juta jiwa) (rata-rata 3 tahun)",
"number of women of reproductive age (15-49 years) affected by anemia (million)":
"Jumlah perempuan usia reproduksi (15-49 tahun) yang menderita anemia (juta jiwa)",
"percentage of children under 5 years affected by wasting (percent)":
"Persentase anak di bawah 5 tahun yang mengalami wasting (persen)",
"prevalence of anemia among women of reproductive age (15-49 years) (percent)":
"Prevalensi anemia pada perempuan usia reproduksi (15-49 tahun) (persen)",
"coefficient of variation of habitual caloric consumption distribution (real number)":
"Koefisien variasi distribusi konsumsi kalori kebiasaan (bilangan riil)",
"incidence of caloric losses at retail distribution level (percent)":
"Insidensi kehilangan kalori pada tingkat distribusi ritel (persen)",
"number of children under 5 years of age who are overweight (modeled estimates) (million)":
"Jumlah anak di bawah 5 tahun yang mengalami kelebihan berat badan (estimasi model) (juta jiwa)",
"number of children under 5 years of age who are stunted (modeled estimates) (million)":
"Jumlah anak di bawah 5 tahun yang mengalami stunting (estimasi model) (juta jiwa)",
"number of newborns with low birthweight (million)":
"Jumlah bayi baru lahir dengan berat badan lahir rendah (juta jiwa)",
"number of obese adults (18 years and older) (million)":
"Jumlah orang dewasa yang mengalami obesitas (18 tahun ke atas) (juta jiwa)",
"percentage of children under 5 years of age who are overweight (modelled estimates) (percent)":
"Persentase anak di bawah 5 tahun yang mengalami kelebihan berat badan (estimasi model) (persen)",
"percentage of children under 5 years of age who are stunted (modelled estimates) (percent)":
"Persentase anak di bawah 5 tahun yang mengalami stunting (estimasi model) (persen)",
"prevalence of low birthweight (percent)":
"Prevalensi berat badan lahir rendah (persen)",
"prevalence of moderate or severe food insecurity in the female adult population (percent) (3-year average)":
"Prevalensi kerawanan pangan sedang atau berat pada penduduk perempuan dewasa (persen) (rata-rata 3 tahun)",
"prevalence of moderate or severe food insecurity in the male adult population (percent) (3-year average)":
"Prevalensi kerawanan pangan sedang atau berat pada penduduk laki-laki dewasa (persen) (rata-rata 3 tahun)",
"prevalence of moderate or severe food insecurity in the total population (percent) (3-year average)":
"Prevalensi kerawanan pangan sedang atau berat pada total penduduk (persen) (rata-rata 3 tahun)",
"prevalence of obesity in the adult population (18 years and older) (percent)":
"Prevalensi obesitas pada penduduk dewasa (18 tahun ke atas) (persen)",
"prevalence of severe food insecurity in the female adult population (percent) (3-year average)":
"Prevalensi kerawanan pangan berat pada penduduk perempuan dewasa (persen) (rata-rata 3 tahun)",
"prevalence of severe food insecurity in the male adult population (percent) (3-year average)":
"Prevalensi kerawanan pangan berat pada penduduk laki-laki dewasa (persen) (rata-rata 3 tahun)",
"prevalence of severe food insecurity in the total population (percent) (3-year average)":
"Prevalensi kerawanan pangan berat pada total penduduk (persen) (rata-rata 3 tahun)",
}
def get_country_name_id(country_name: str) -> str:
return COUNTRY_NAME_ID_MAP.get(str(country_name).strip(), str(country_name))
def get_indicator_name_id(indicator_name: str) -> str:
return INDICATOR_NAME_ID_MAP.get(str(indicator_name).lower().strip(), str(indicator_name))
def get_pillar_name_id(pillar_name: str) -> str:
return PILLAR_NAME_ID_MAP.get(str(pillar_name).strip(), str(pillar_name))
# =============================================================================
# SDG-ONLY KEYWORD SET
# =============================================================================
SDG_ONLY_KEYWORDS: frozenset = frozenset([
"prevalence of undernourishment (percent) (3-year average)",
"number of people undernourished (million) (3-year average)",
"prevalence of severe food insecurity in the total population (percent) (3-year average)",
"prevalence of severe food insecurity in the male adult population (percent) (3-year average)",
"prevalence of severe food insecurity in the female adult population (percent) (3-year average)",
"prevalence of moderate or severe food insecurity in the total population (percent) (3-year average)",
"prevalence of moderate or severe food insecurity in the male adult population (percent) (3-year average)",
"prevalence of moderate or severe food insecurity in the female adult population (percent) (3-year average)",
"number of severely food insecure people (million) (3-year average)",
"number of severely food insecure male adults (million) (3-year average)",
"number of severely food insecure female adults (million) (3-year average)",
"number of moderately or severely food insecure people (million) (3-year average)",
"number of moderately or severely food insecure male adults (million) (3-year average)",
"number of moderately or severely food insecure female adults (million) (3-year average)",
"percentage of children under 5 years of age who are stunted (modelled estimates) (percent)",
"number of children under 5 years of age who are stunted (modeled estimates) (million)",
"percentage of children under 5 years affected by wasting (percent)",
"number of children under 5 years affected by wasting (million)",
"percentage of children under 5 years of age who are overweight (modelled estimates) (percent)",
"number of children under 5 years of age who are overweight (modeled estimates) (million)",
"prevalence of anemia among women of reproductive age (15-49 years) (percent)",
"number of women of reproductive age (15-49 years) affected by anemia (million)",
])
_SDG_ONLY_LOWER: frozenset = frozenset(k.lower() for k in SDG_ONLY_KEYWORDS)
_FIES_DETECTION_KEYWORDS: frozenset = frozenset([
"prevalence of severe food insecurity in the total population (percent) (3-year average)",
"prevalence of moderate or severe food insecurity in the total population (percent) (3-year average)",
"number of severely food insecure people (million) (3-year average)",
"number of moderately or severely food insecure people (million) (3-year average)",
])
_FIES_DETECTION_LOWER: frozenset = frozenset(k.lower() for k in _FIES_DETECTION_KEYWORDS)
DIRECTION_INVERT_KEYWORDS = frozenset({
"negative", "lower_better", "lower_is_better", "inverse", "neg",
})
DIRECTION_POSITIVE_KEYWORDS = frozenset({
"positive", "higher_better", "higher_is_better",
})
_PERFORMANCE_THRESHOLD: float = 60.0
# =============================================================================
# PURE HELPERS
# =============================================================================
def _should_invert(direction: str, logger=None, context: str = "") -> bool:
d = str(direction).lower().strip()
if d in DIRECTION_INVERT_KEYWORDS:
return True
if d in DIRECTION_POSITIVE_KEYWORDS:
return False
if logger:
logger.warning(
f" [DIRECTION WARNING] Unknown direction '{direction}' "
f"{'(' + context + ')' if context else ''}. Defaulting to positive (no invert)."
)
return False
def global_minmax(series: pd.Series, lo: float = 1.0, hi: float = 100.0) -> pd.Series:
values = series.dropna().values
if len(values) == 0:
return pd.Series(np.nan, index=series.index)
v_min, v_max = values.min(), values.max()
if v_min == v_max:
return pd.Series((lo + hi) / 2.0, index=series.index)
result = np.full(len(series), np.nan)
not_nan = series.notna()
result[not_nan.values] = lo + (series[not_nan].values - v_min) / (v_max - v_min) * (hi - lo)
return pd.Series(result, index=series.index)
def _compute_yoy(df: pd.DataFrame) -> pd.DataFrame:
df = df.sort_values("year").copy()
df["value_prev"] = df["value"].shift(1)
df["norm_value_prev"] = df["norm_value"].shift(1)
df["yoy_value"] = np.where(
df["value"].notna() & df["value_prev"].notna(),
df["value"] - df["value_prev"],
np.nan,
)
df["yoy_norm_value"] = np.where(
df["norm_value"].notna() & df["norm_value_prev"].notna(),
df["norm_value"] - df["norm_value_prev"],
np.nan,
)
df = df.drop(columns=["value_prev", "norm_value_prev"])
return df
def _is_lower_better(direction: str) -> bool:
return str(direction).lower().strip() in DIRECTION_INVERT_KEYWORDS
# =============================================================================
# NARRATIVE CONDITION DETECTORS
# =============================================================================
def _detect_trend(scores_by_year: pd.Series, lower_better: bool) -> str:
if len(scores_by_year) < 3:
return "insufficient_data"
years = sorted(scores_by_year.index)
vals = [scores_by_year[y] for y in years if not pd.isna(scores_by_year.get(y, np.nan))]
if len(vals) < 3:
return "insufficient_data"
x = np.arange(len(vals))
slope = np.polyfit(x, vals, 1)[0]
improving = (slope > 0 and not lower_better) or (slope < 0 and lower_better)
mid = len(vals) // 2
first_half = vals[:mid]
second_half = vals[mid:]
slope1 = np.polyfit(np.arange(len(first_half)), first_half, 1)[0] if len(first_half) > 1 else 0
slope2 = np.polyfit(np.arange(len(second_half)), second_half, 1)[0] if len(second_half) > 1 else 0
cv = np.std(vals) / (np.mean(vals) + 1e-9)
if cv > 0.25:
return "fluctuating"
if improving:
if lower_better:
slowing = slope2 > slope1
else:
slowing = slope2 < slope1
return "improving_slowing" if slowing else "improving_consistent"
else:
return "deteriorating"
def _detect_gap_trend(df_ind: pd.DataFrame, lower_better: bool) -> str:
# Hanya hitung gap di antara negara asli (bukan ASEAN)
df_real = df_ind[df_ind["country_id"] != ASEAN_COUNTRY_ID]
std_by_year = (
df_real.groupby("year")["value"]
.std()
.dropna()
)
if len(std_by_year) < 3:
return "unknown"
years = sorted(std_by_year.index)
stds = [std_by_year[y] for y in years]
slope = np.polyfit(np.arange(len(stds)), stds, 1)[0]
if abs(slope) < 0.01 * np.mean(stds):
return "stable"
return "widening" if slope > 0 else "narrowing"
def _detect_anomaly_year(scores_by_year: pd.Series) -> tuple:
if len(scores_by_year) < 3:
return None, None
years = sorted(scores_by_year.index)
deltas = {}
for i in range(1, len(years)):
y_prev = years[i - 1]
y_curr = years[i]
v_prev = scores_by_year.get(y_prev, np.nan)
v_curr = scores_by_year.get(y_curr, np.nan)
if not pd.isna(v_prev) and not pd.isna(v_curr):
deltas[y_curr] = v_curr - v_prev
if not deltas:
return None, None
max_drop_year = min(deltas, key=deltas.get)
max_rise_year = max(deltas, key=deltas.get)
threshold = 1.5 * np.std(list(deltas.values()))
if abs(deltas[max_drop_year]) > threshold and deltas[max_drop_year] < 0:
return max_drop_year, "drop"
if abs(deltas[max_rise_year]) > threshold and deltas[max_rise_year] > 0:
return max_rise_year, "rise"
return None, None
def _detect_consistency(df_ind: pd.DataFrame, lower_better: bool) -> tuple:
"""Hanya negara asli (bukan ASEAN) yang di-analisa konsistensinya."""
df_real = df_ind[df_ind["country_id"] != ASEAN_COUNTRY_ID]
country_avg = (
df_real.groupby("country_name")["value"]
.mean()
.dropna()
)
if country_avg.empty:
return None, None, False
if lower_better:
best = country_avg.idxmin()
worst = country_avg.idxmax()
else:
best = country_avg.idxmax()
worst = country_avg.idxmin()
asean_avg_by_year = df_real.groupby("year")["value"].mean()
country_by_year = df_real[df_real["country_name"] == best].set_index("year")["value"]
years_both = set(asean_avg_by_year.index) & set(country_by_year.index)
if not years_both:
return best, worst, False
if lower_better:
consistent = all(
country_by_year[y] <= asean_avg_by_year[y]
for y in years_both
if not pd.isna(country_by_year.get(y, np.nan))
)
else:
consistent = all(
country_by_year[y] >= asean_avg_by_year[y]
for y in years_both
if not pd.isna(country_by_year.get(y, np.nan))
)
return best, worst, consistent
# =============================================================================
# NARRATIVE BUILDER — PER INDICATOR PER YEAR (1 pillar, 1 tahun)
# =============================================================================
#
# Granularity agg_narrative_indicator SEKARANG per (indicator_id, year), bukan
# lagi per indicator_id gabungan seluruh tahun. Setiap baris narasi menjelaskan
# posisi 1 indikator dalam 1 pillar, pada 1 tahun tertentu: nilai regional,
# skor ternormalisasi, peringkat indikator tsb di antara indikator lain dalam
# pillar yang sama pada tahun yang sama, perubahan YoY, serta negara terbaik/
# terlemah untuk indikator tsb pada tahun tersebut.
# =============================================================================
def _build_narrative_per_indicator_year(row: pd.Series) -> tuple:
"""
row diharapkan datang dari baris ASEAN (country_id=0) pada agg_indicator_norm,
sudah digabung dengan kolom rank_in_pillar_year, n_indicators_in_pillar_year,
country_best, country_worst (dihitung dari negara asli untuk indicator+year yang sama).
"""
ind_name_en = str(row["indicator_name"]).strip()
ind_name_id = str(row.get("indicator_name_id", ind_name_en)).strip()
unit = str(row["unit"]).strip() if row["unit"] else ""
pillar_en = str(row["pillar_name"]).strip()
pillar_id_ = get_pillar_name_id(pillar_en)
framework = str(row["framework"]).strip()
year = int(row["year"])
value = row.get("value", np.nan)
score = row.get("norm_score_1_100", np.nan)
yoy = row.get("yoy_value", np.nan)
rank = row.get("rank_in_pillar_year", np.nan)
n_ind = row.get("n_indicators_in_pillar_year", np.nan)
best_country_en = row.get("country_best")
worst_country_en = row.get("country_worst")
best_country_id = row.get("country_best_id")
worst_country_id = row.get("country_worst_id")
def fmt(v):
if pd.isna(v):
return "N/A"
abs_v = abs(v)
s = f"{v:,.1f}" if abs_v >= 1000 else (f"{v:.2f}" if abs_v >= 10 else f"{v:.3f}")
return f"{s} {unit}".strip() if unit else s
sentences_en = []
sentences_id = []
perf_word_en = "good" if (pd.notna(score) and score >= _PERFORMANCE_THRESHOLD) else "below target"
perf_word_id = "baik" if (pd.notna(score) and score >= _PERFORMANCE_THRESHOLD) else "di bawah target"
if pd.notna(rank) and pd.notna(n_ind) and int(n_ind) > 0:
rank_i = int(rank)
n_i = int(n_ind)
suffix = {1: "st", 2: "nd", 3: "rd"}.get(rank_i, "th")
s1_en = (
f"In {year}, {ind_name_en} ({framework}, {pillar_en}) recorded a regional average of "
f"{fmt(value)}, ranking {rank_i}{suffix} out of {n_i} indicators within the {pillar_en} "
f"pillar for that year, with a normalized score of {fmt(score)} ({perf_word_en})."
)
s1_id = (
f"Pada tahun {year}, {ind_name_id} ({framework}, {pillar_id_}) mencatat rata-rata "
f"regional sebesar {fmt(value)}, menempati peringkat {rank_i} dari {n_i} indikator "
f"dalam pilar {pillar_id_} pada tahun tersebut, dengan skor ternormalisasi "
f"{fmt(score)} ({perf_word_id})."
)
else:
s1_en = (
f"{ind_name_en} ({framework}, {pillar_en}, {year}): regional average was {fmt(value)}, "
f"with a normalized score of {fmt(score)} ({perf_word_en})."
)
s1_id = (
f"{ind_name_id} ({framework}, {pillar_id_}, {year}): rata-rata regional sebesar "
f"{fmt(value)}, dengan skor ternormalisasi {fmt(score)} ({perf_word_id})."
)
sentences_en.append(s1_en)
sentences_id.append(s1_id)
if pd.notna(yoy):
if abs(yoy) < 1e-9:
s2_en = "This value was unchanged compared to the previous year."
s2_id = "Nilai ini tidak berubah dibandingkan tahun sebelumnya."
elif yoy > 0:
s2_en = f"This represents an increase of {fmt(abs(yoy))} from the previous year."
s2_id = f"Ini menunjukkan kenaikan sebesar {fmt(abs(yoy))} dari tahun sebelumnya."
else:
s2_en = f"This represents a decrease of {fmt(abs(yoy))} from the previous year."
s2_id = f"Ini menunjukkan penurunan sebesar {fmt(abs(yoy))} dari tahun sebelumnya."
sentences_en.append(s2_en)
sentences_id.append(s2_id)
else:
s2_en = "Year-over-year comparison is not available for this year."
s2_id = "Perbandingan tahun-ke-tahun tidak tersedia untuk tahun ini."
sentences_en.append(s2_en)
sentences_id.append(s2_id)
if best_country_en and worst_country_en and best_country_en != worst_country_en:
s3_en = (
f"Among ASEAN countries in {year}, {best_country_en} recorded the best performance "
f"for this indicator, while {worst_country_en} recorded the weakest."
)
s3_id = (
f"Di antara negara ASEAN pada tahun {year}, {best_country_id} mencatat performa "
f"terbaik untuk indikator ini, sementara {worst_country_id} mencatat performa terlemah."
)
sentences_en.append(s3_en)
sentences_id.append(s3_id)
narrative_en = " ".join(s for s in sentences_en if s)
narrative_id = " ".join(s for s in sentences_id if s)
return narrative_en, narrative_id
# =============================================================================
# MAIN CLASS
# =============================================================================
class IndicatorNormAggregator:
def __init__(self, client: bigquery.Client):
self.client = client
self.logger = logging.getLogger(self.__class__.__name__)
self.logger.propagate = False
self.df = None
self.df_unit = None
self.sdgs_start_year = None
self.pipeline_start = None
self.pipeline_metadata = {
"rows_fetched" : 0,
"rows_loaded" : 0,
"rows_loaded_narrative" : 0,
"start_time" : None,
"end_time" : None,
}
# =========================================================================
# STEP 1: Load data
# =========================================================================
def load_data(self):
self.logger.info("\n" + "=" * 80)
self.logger.info("STEP 1: LOAD DATA — fact_asean_food_security_selected")
self.logger.info("=" * 80)
self.df = read_from_bigquery(
self.client, "fact_asean_food_security_selected", layer="gold"
)
required = {
"country_id", "country_name",
"indicator_id", "indicator_name", "direction",
"pillar_id", "pillar_name",
"year", "value",
}
missing = required - set(self.df.columns)
if missing:
raise ValueError(f"Kolom tidak ditemukan: {missing}")
n_null = self.df["direction"].isna().sum()
if n_null > 0:
self.logger.warning(f" {n_null} rows direction NULL -> diisi 'positive'")
self.df["direction"] = self.df["direction"].fillna("positive")
# Rename pillar_name: add 'Food ' prefix, remove
PILLAR_RENAME_MAP = {
'Availability' : 'Food Availability',
'Access' : 'Food Access',
'Utilization' : 'Food Utilization',
'Stability' : 'Food Stability',
'Other' : 'Food Other',
}
self.df["pillar_name"] = self.df["pillar_name"].replace(PILLAR_RENAME_MAP)
self.pipeline_metadata["rows_fetched"] = len(self.df)
self.logger.info(f" Rows : {len(self.df):,}")
self.logger.info(f" Countries : {self.df['country_id'].nunique()}")
self.logger.info(f" Indicators: {self.df['indicator_id'].nunique()}")
self.logger.info(
f" Years : {int(self.df['year'].min())} - {int(self.df['year'].max())}"
)
# =========================================================================
# STEP 2: Load unit
# =========================================================================
def load_units(self):
self.logger.info("\n" + "=" * 80)
self.logger.info("STEP 2: LOAD UNIT — dim_indicator")
self.logger.info("=" * 80)
dim = read_from_bigquery(self.client, "dim_indicator", layer="gold")
if "indicator_id" not in dim.columns or "unit" not in dim.columns:
raise ValueError(
f"dim_indicator harus punya kolom 'indicator_id' dan 'unit'. "
f"Kolom tersedia: {list(dim.columns)}"
)
self.df_unit = (
dim[["indicator_id", "unit"]]
.drop_duplicates(subset=["indicator_id"])
.copy()
)
self.df_unit["indicator_id"] = self.df_unit["indicator_id"].astype(int)
self.df_unit["unit"] = self.df_unit["unit"].fillna("").astype(str)
self.logger.info(f" dim_indicator rows (unique indicator_id): {len(self.df_unit):,}")
# =========================================================================
# STEP 3: Merge unit
# =========================================================================
def _merge_unit(self):
before = len(self.df)
self.df = self.df.merge(self.df_unit, on="indicator_id", how="left")
self.df["unit"] = self.df["unit"].fillna("").astype(str)
after = len(self.df)
assert before == after, f"Row count berubah: {before} -> {after}"
self.logger.info(f" Merge unit OK. Rows: {after:,}")
# =========================================================================
# STEP 3b: Tambah kolom nama Bahasa Indonesia
# =========================================================================
def _add_indonesia_name_columns(self):
self.df["country_name_id"] = self.df["country_name"].apply(get_country_name_id).astype(str)
self.df["indicator_name_id"] = self.df["indicator_name"].apply(get_indicator_name_id).astype(str)
self.df["pillar_name_id"] = self.df["pillar_name"].apply(get_pillar_name_id).astype(str)
self.logger.info(" Kolom terjemahan Indonesia ditambahkan.")
sample_pil = self.df[["pillar_name", "pillar_name_id"]].drop_duplicates()
self.logger.info(" Pillar mapping:")
for _, r in sample_pil.iterrows():
self.logger.info(f" {r['pillar_name']:<20} -> {r['pillar_name_id']}")
# =========================================================================
# STEP 4: Deteksi sdgs_start_year
# =========================================================================
def _detect_sdgs_start_year(self) -> int:
fies_rows = self.df[
self.df["indicator_name"].str.lower().str.strip().isin(_FIES_DETECTION_LOWER)
]
if not fies_rows.empty:
sdgs_start = int(fies_rows["year"].min())
self.logger.info(f" [Metode 1 - FIES explicit] sdgs_start_year = {sdgs_start}")
return sdgs_start
ind_min_year = (
self.df.groupby("indicator_id")["year"]
.min().reset_index()
.rename(columns={"year": "min_year"})
)
unique_years = sorted(ind_min_year["min_year"].unique())
if len(unique_years) == 1:
sdgs_start = int(unique_years[0]) + 9999
else:
gaps = [
(unique_years[i+1] - unique_years[i], unique_years[i], unique_years[i+1])
for i in range(len(unique_years) - 1)
]
gaps.sort(reverse=True)
_, y_before, y_after = gaps[0]
sdgs_start = int(y_after)
self.logger.info(f" Gap terbesar: {y_before} -> {y_after} -> sdgs_start_year = {sdgs_start}")
return sdgs_start
# =========================================================================
# STEP 5: Assign framework
# =========================================================================
def _assign_framework(self):
df = self.df.copy()
df["_is_sdg_kw"] = df["indicator_name"].str.lower().str.strip().isin(_SDG_ONLY_LOWER)
df["framework"] = "MDGs"
mask_sdgs = df["_is_sdg_kw"] & (df["year"] >= self.sdgs_start_year)
df.loc[mask_sdgs, "framework"] = "SDGs"
df = df.drop(columns=["_is_sdg_kw"])
self.df = df
fw_dist = self.df["framework"].value_counts()
self.logger.info("\n Framework distribution (rows):")
for fw, cnt in fw_dist.items():
self.logger.info(f" {fw:<6}: {cnt:,} rows")
# =========================================================================
# STEP 6: Hitung norm_value per indikator
# =========================================================================
def _compute_norm_values(self) -> pd.DataFrame:
df = self.df.copy()
norm_parts = []
for ind_id, grp in df.groupby("indicator_id"):
grp = grp.copy()
direction = str(grp["direction"].iloc[0])
do_invert = _should_invert(
direction, self.logger, context=f"indicator_id={ind_id}"
)
valid_mask = grp["value"].notna()
n_valid = valid_mask.sum()
if n_valid < 2:
grp["norm_value"] = np.nan
norm_parts.append(grp)
continue
raw = grp.loc[valid_mask, "value"].values
v_min = raw.min()
v_max = raw.max()
normed = np.full(len(grp), np.nan)
if v_min == v_max:
normed[valid_mask.values] = 0.5
else:
normed[valid_mask.values] = (raw - v_min) / (v_max - v_min)
if do_invert:
normed = np.where(np.isnan(normed), np.nan, 1.0 - normed)
grp["norm_value"] = normed
norm_parts.append(grp)
return pd.concat(norm_parts, ignore_index=True)
# =========================================================================
# STEP 6b: Tambah baris ASEAN (rata-rata dari semua negara per ind per year)
# =========================================================================
def _add_asean_rows(self, df_normed: pd.DataFrame) -> pd.DataFrame:
"""
Buat baris ASEAN aggregate: nilai = rata-rata nilai semua negara asli.
norm_value = rata-rata norm_value semua negara asli.
"""
real_rows = df_normed[df_normed["country_id"] != ASEAN_COUNTRY_ID]
dim_cols = [
"indicator_id", "indicator_name", "indicator_name_id",
"unit", "direction",
"pillar_id", "pillar_name", "pillar_name_id",
"framework",
]
asean_agg = (
real_rows.groupby(["indicator_id", "year"])
.agg(
value =("value", "mean"),
norm_value=("norm_value", "mean"),
)
.reset_index()
)
# Gabung kolom dimensi dari baris pertama per indicator_id
dim_ref = (
real_rows[dim_cols]
.drop_duplicates(subset=["indicator_id"])
.copy()
)
asean_agg = asean_agg.merge(dim_ref, on="indicator_id", how="left")
asean_agg["country_id"] = ASEAN_COUNTRY_ID
asean_agg["country_name"] = ASEAN_COUNTRY_NAME
asean_agg["country_name_id"] = ASEAN_COUNTRY_NAME_ID
# Hanya kolom yang ada di df_normed
cols_needed = [c for c in df_normed.columns if c in asean_agg.columns]
for c in df_normed.columns:
if c not in asean_agg.columns:
asean_agg[c] = np.nan
return pd.concat([df_normed, asean_agg[df_normed.columns]], ignore_index=True)
# =========================================================================
# STEP 7: Hitung YoY
# =========================================================================
def _compute_yoy_columns(self, df: pd.DataFrame) -> pd.DataFrame:
parts = []
groups = df.groupby(["indicator_id", "country_id"], sort=False)
for (ind_id, country_id), grp in groups:
parts.append(_compute_yoy(grp))
return pd.concat(parts, ignore_index=True)
# =========================================================================
# STEP 8: Scale ke 1-100
# =========================================================================
def _compute_scores(self, df: pd.DataFrame) -> pd.DataFrame:
score_parts = []
for ind_id, grp in df.groupby("indicator_id"):
grp = grp.copy()
grp["norm_score_1_100"] = global_minmax(grp["norm_value"])
score_parts.append(grp)
return pd.concat(score_parts, ignore_index=True)
# =========================================================================
# STEP 9: Assign performance label
# =========================================================================
def _assign_performance(self, df: pd.DataFrame) -> pd.DataFrame:
df = df.copy()
df["performance"] = pd.NA
has_score = df["norm_score_1_100"].notna()
df.loc[has_score & (df["norm_score_1_100"] >= _PERFORMANCE_THRESHOLD), "performance"] = "Good"
df.loc[has_score & (df["norm_score_1_100"] < _PERFORMANCE_THRESHOLD), "performance"] = "Bad"
return df
# =========================================================================
# STEP 10: Save agg_indicator_norm (termasuk ASEAN rows)
# =========================================================================
def _save(self, df: pd.DataFrame) -> int:
table_name = "agg_indicator_norm"
out = df[[
"year", "country_id", "country_name", "country_name_id",
"indicator_id", "indicator_name", "indicator_name_id",
"unit", "direction",
"pillar_id", "pillar_name", "pillar_name_id",
"framework",
"value", "norm_value", "norm_score_1_100",
"yoy_value", "yoy_norm_value", "performance",
]].copy()
out = out.sort_values(
["year", "country_name", "pillar_name", "indicator_name"]
).reset_index(drop=True)
out["year"] = out["year"].astype(int)
out["country_id"] = out["country_id"].astype(int)
out["country_name"] = out["country_name"].astype(str)
out["country_name_id"] = out["country_name_id"].astype(str)
out["indicator_id"] = out["indicator_id"].astype(int)
out["indicator_name"] = out["indicator_name"].astype(str)
out["indicator_name_id"] = out["indicator_name_id"].astype(str)
out["unit"] = out["unit"].astype(str)
out["direction"] = out["direction"].astype(str)
out["pillar_id"] = out["pillar_id"].astype(int)
out["pillar_name"] = out["pillar_name"].astype(str)
out["pillar_name_id"] = out["pillar_name_id"].astype(str)
out["framework"] = out["framework"].astype(str)
out["value"] = out["value"].astype(float)
out["norm_value"] = out["norm_value"].astype(float)
out["norm_score_1_100"] = out["norm_score_1_100"].astype(float)
out["yoy_value"] = pd.to_numeric(out["yoy_value"], errors="coerce").astype(float)
out["yoy_norm_value"] = pd.to_numeric(out["yoy_norm_value"], errors="coerce").astype(float)
out["performance"] = out["performance"].astype(str).replace("nan", pd.NA).astype("string")
n_asean = (out["country_id"] == ASEAN_COUNTRY_ID).sum()
n_country = (out["country_id"] != ASEAN_COUNTRY_ID).sum()
self.logger.info(f" Total rows : {len(out):,} ({n_country:,} country + {n_asean:,} ASEAN)")
schema = [
bigquery.SchemaField("year", "INTEGER", mode="REQUIRED"),
bigquery.SchemaField("country_id", "INTEGER", mode="REQUIRED"),
bigquery.SchemaField("country_name", "STRING", mode="REQUIRED"),
bigquery.SchemaField("country_name_id", "STRING", mode="NULLABLE"),
bigquery.SchemaField("indicator_id", "INTEGER", mode="REQUIRED"),
bigquery.SchemaField("indicator_name", "STRING", mode="REQUIRED"),
bigquery.SchemaField("indicator_name_id", "STRING", mode="NULLABLE"),
bigquery.SchemaField("unit", "STRING", mode="NULLABLE"),
bigquery.SchemaField("direction", "STRING", mode="REQUIRED"),
bigquery.SchemaField("pillar_id", "INTEGER", mode="REQUIRED"),
bigquery.SchemaField("pillar_name", "STRING", mode="REQUIRED"),
bigquery.SchemaField("pillar_name_id", "STRING", mode="NULLABLE"),
bigquery.SchemaField("framework", "STRING", mode="REQUIRED"),
bigquery.SchemaField("value", "FLOAT", mode="REQUIRED"),
bigquery.SchemaField("norm_value", "FLOAT", mode="NULLABLE"),
bigquery.SchemaField("norm_score_1_100", "FLOAT", mode="NULLABLE"),
bigquery.SchemaField("yoy_value", "FLOAT", mode="NULLABLE"),
bigquery.SchemaField("yoy_norm_value", "FLOAT", mode="NULLABLE"),
bigquery.SchemaField("performance", "STRING", mode="NULLABLE"),
]
rows_loaded = load_to_bigquery(
self.client, out, table_name,
layer="gold", write_disposition="WRITE_TRUNCATE", schema=schema,
)
log_update(self.client, "DW", table_name, "full_load", rows_loaded)
self.logger.info(f" [OK] {table_name}: {rows_loaded:,} rows -> [Gold] fs_asean_gold")
metadata = {
"source_class" : self.__class__.__name__,
"table_name" : table_name,
"execution_timestamp": self.pipeline_start,
"duration_seconds" : (datetime.now() - self.pipeline_start).total_seconds(),
"rows_fetched" : self.pipeline_metadata["rows_fetched"],
"rows_transformed" : rows_loaded,
"rows_loaded" : rows_loaded,
"completeness_pct" : 100.0,
"config_snapshot" : json.dumps({
"sdgs_start_year" : self.sdgs_start_year,
"layer" : "gold",
"normalization" : "per_indicator_global_minmax",
"performance_threshold": _PERFORMANCE_THRESHOLD,
"asean_country_id" : ASEAN_COUNTRY_ID,
"architecture" : "ASEAN merged into country table (country_id=0)",
"pillar_change" : "Renamed to Food Other; all pillars use 'Food ' prefix",
}),
"validation_metrics" : json.dumps({
"total_rows" : rows_loaded,
"n_indicators" : int(out["indicator_id"].nunique()),
"n_countries" : int(out[out["country_id"] != ASEAN_COUNTRY_ID]["country_id"].nunique()),
"asean_rows" : int(n_asean),
}),
}
save_etl_metadata(self.client, metadata)
return rows_loaded
# =========================================================================
# STEP 11: agg_narrative_indicator (per indicator_id PER YEAR, 1 pillar 1 tahun)
# =========================================================================
def _build_narrative_table(self, df_final: pd.DataFrame):
self.logger.info("\n" + "=" * 80)
self.logger.info("STEP 11: agg_narrative_indicator")
self.logger.info(" Granularity: per indicator_id PER YEAR (1 pillar, 1 tahun)")
self.logger.info(" Narasi menjelaskan posisi indikator dalam pillar-nya, per tahun")
self.logger.info("=" * 80)
# Negara asli saja untuk analisa negara terbaik/terlemah per tahun
df_real = df_final[df_final["country_id"] != ASEAN_COUNTRY_ID].copy()
# Baris ASEAN (regional average) menjadi basis nilai per indicator per year
df_asean = df_final[df_final["country_id"] == ASEAN_COUNTRY_ID].copy()
if df_asean.empty:
self.logger.warning(" [WARNING] Tidak ada baris ASEAN; agg_narrative_indicator kosong.")
df_asean = df_final.copy()
# ---- Rank indikator dalam pillar yang sama, pada tahun yang sama ----
# (dihitung dari norm_score_1_100 baris ASEAN/regional; skor sudah
# searah -- semakin tinggi semakin baik -- untuk semua indikator)
df_asean["rank_in_pillar_year"] = (
df_asean.groupby(["pillar_id", "year"])["norm_score_1_100"]
.rank(method="min", ascending=False)
)
df_asean["n_indicators_in_pillar_year"] = (
df_asean.groupby(["pillar_id", "year"])["indicator_id"]
.transform("nunique")
)
# ---- Negara terbaik / terlemah per indikator, per tahun (negara asli) ----
def _best_worst(g: pd.DataFrame) -> pd.Series:
g_valid = g[g["norm_score_1_100"].notna()]
if g_valid.empty:
return pd.Series({"country_best": None, "country_worst": None})
best_row = g_valid.loc[g_valid["norm_score_1_100"].idxmax()]
worst_row = g_valid.loc[g_valid["norm_score_1_100"].idxmin()]
return pd.Series({
"country_best" : best_row["country_name"],
"country_worst": worst_row["country_name"],
})
country_stats = (
df_real.groupby(["indicator_id", "year"])
.apply(_best_worst)
.reset_index()
)
country_stats["country_best_id"] = country_stats["country_best"].apply(
lambda x: get_country_name_id(x) if pd.notna(x) and x is not None else None
)
country_stats["country_worst_id"] = country_stats["country_worst"].apply(
lambda x: get_country_name_id(x) if pd.notna(x) and x is not None else None
)
# ---- Gabung ----
df_agg = df_asean.merge(country_stats, on=["indicator_id", "year"], how="left")
# ---- Build narrative per baris (per indicator_id per year) ----
narratives_en = []
narratives_id = []
for _, row in df_agg.iterrows():
n_en, n_id = _build_narrative_per_indicator_year(row)
narratives_en.append(n_en)
narratives_id.append(n_id)
df_agg["narrative_en"] = narratives_en
df_agg["narrative_id"] = narratives_id
# ---- Output ----
out = df_agg[[
"year",
"indicator_id", "indicator_name", "indicator_name_id",
"unit", "direction",
"pillar_id", "pillar_name", "pillar_name_id",
"framework",
"value", "norm_score_1_100", "performance",
"yoy_value",
"rank_in_pillar_year", "n_indicators_in_pillar_year",
"country_best", "country_worst",
"country_best_id", "country_worst_id",
"narrative_en", "narrative_id",
]].copy()
out = out.sort_values(["year", "pillar_name", "rank_in_pillar_year", "indicator_name"]).reset_index(drop=True)
out["year"] = out["year"].astype(int)
out["indicator_id"] = out["indicator_id"].astype(int)
out["indicator_name"] = out["indicator_name"].astype(str)
out["indicator_name_id"] = out["indicator_name_id"].astype(str)
out["unit"] = out["unit"].fillna("").astype(str)
out["direction"] = out["direction"].astype(str)
out["pillar_id"] = out["pillar_id"].astype(int)
out["pillar_name"] = out["pillar_name"].astype(str)
out["pillar_name_id"] = out["pillar_name_id"].astype(str)
out["framework"] = out["framework"].astype(str)
out["value"] = pd.to_numeric(out["value"], errors="coerce").astype(float)
out["norm_score_1_100"] = pd.to_numeric(out["norm_score_1_100"], errors="coerce").astype(float)
out["performance"] = out["performance"].astype(str).replace("nan", pd.NA).astype("string")
out["yoy_value"] = pd.to_numeric(out["yoy_value"], errors="coerce").astype(float)
out["rank_in_pillar_year"] = pd.to_numeric(out["rank_in_pillar_year"], errors="coerce").astype("Int64")
out["n_indicators_in_pillar_year"] = pd.to_numeric(out["n_indicators_in_pillar_year"], errors="coerce").astype("Int64")
out["country_best"] = out["country_best"].astype(str).replace("nan", pd.NA).astype("string")
out["country_worst"] = out["country_worst"].astype(str).replace("nan", pd.NA).astype("string")
out["country_best_id"] = out["country_best_id"].astype(str).replace("nan", pd.NA).astype("string")
out["country_worst_id"] = out["country_worst_id"].astype(str).replace("nan", pd.NA).astype("string")
out["narrative_en"] = out["narrative_en"].astype(str)
out["narrative_id"] = out["narrative_id"].astype(str)
schema = [
bigquery.SchemaField("year", "INTEGER", mode="REQUIRED"),
bigquery.SchemaField("indicator_id", "INTEGER", mode="REQUIRED"),
bigquery.SchemaField("indicator_name", "STRING", mode="REQUIRED"),
bigquery.SchemaField("indicator_name_id", "STRING", mode="NULLABLE"),
bigquery.SchemaField("unit", "STRING", mode="NULLABLE"),
bigquery.SchemaField("direction", "STRING", mode="REQUIRED"),
bigquery.SchemaField("pillar_id", "INTEGER", mode="REQUIRED"),
bigquery.SchemaField("pillar_name", "STRING", mode="REQUIRED"),
bigquery.SchemaField("pillar_name_id", "STRING", mode="NULLABLE"),
bigquery.SchemaField("framework", "STRING", mode="REQUIRED"),
bigquery.SchemaField("value", "FLOAT", mode="NULLABLE"),
bigquery.SchemaField("norm_score_1_100", "FLOAT", mode="NULLABLE"),
bigquery.SchemaField("performance", "STRING", mode="NULLABLE"),
bigquery.SchemaField("yoy_value", "FLOAT", mode="NULLABLE"),
bigquery.SchemaField("rank_in_pillar_year", "INTEGER", mode="NULLABLE"),
bigquery.SchemaField("n_indicators_in_pillar_year", "INTEGER", mode="NULLABLE"),
bigquery.SchemaField("country_best", "STRING", mode="NULLABLE"),
bigquery.SchemaField("country_worst", "STRING", mode="NULLABLE"),
bigquery.SchemaField("country_best_id", "STRING", mode="NULLABLE"),
bigquery.SchemaField("country_worst_id", "STRING", mode="NULLABLE"),
bigquery.SchemaField("narrative_en", "STRING", mode="NULLABLE"),
bigquery.SchemaField("narrative_id", "STRING", mode="NULLABLE"),
]
rows_loaded = load_to_bigquery(
self.client, out, "agg_narrative_indicator",
layer="gold", write_disposition="WRITE_TRUNCATE", schema=schema,
)
log_update(self.client, "DW", "agg_narrative_indicator", "full_load", rows_loaded)
self.logger.info(
f" [OK] agg_narrative_indicator: {rows_loaded:,} rows -> [Gold] fs_asean_gold"
)
metadata = {
"source_class" : self.__class__.__name__,
"table_name" : "agg_narrative_indicator",
"execution_timestamp": self.pipeline_start,
"duration_seconds" : (datetime.now() - self.pipeline_start).total_seconds(),
"rows_fetched" : self.pipeline_metadata["rows_fetched"],
"rows_transformed" : rows_loaded,
"rows_loaded" : rows_loaded,
"completeness_pct" : 100.0,
"config_snapshot" : json.dumps({
"granularity" : "indicator_id x year (1 pillar, 1 tahun)",
"narrative_style" : "interpretive, plain text, bilingual EN/ID",
"rank_basis" : "norm_score_1_100 within same pillar_id and year",
"architecture" : "Built from ASEAN rows (country_id=0) in agg_indicator_norm",
"pillar_change" : "Renamed to Food Other; all pillars use 'Food ' prefix",
}),
"validation_metrics" : json.dumps({
"total_rows" : rows_loaded,
"n_indicators": int(out["indicator_id"].nunique()),
"n_years" : int(out["year"].nunique()),
}),
}
save_etl_metadata(self.client, metadata)
self.pipeline_metadata["rows_loaded_narrative"] = rows_loaded
# =========================================================================
# RUN
# =========================================================================
def run(self):
self.pipeline_start = datetime.now()
self.pipeline_metadata["start_time"] = self.pipeline_start
self.logger.info("\n" + "=" * 80)
self.logger.info("INDICATOR NORM AGGREGATION")
self.logger.info(" ASEAN rows ditambahkan ke agg_indicator_norm (country_id=0)")
self.logger.info(" Rename Food Other; all pillars use Food prefix")
self.logger.info(" agg_narrative_indicator: granularity per indicator per year (1 pillar, 1 tahun)")
self.logger.info("=" * 80)
self.load_data()
self.load_units()
self._merge_unit()
self._add_indonesia_name_columns()
self.sdgs_start_year = self._detect_sdgs_start_year()
self._assign_framework()
df_normed = self._compute_norm_values()
df_with_asean = self._add_asean_rows(df_normed) # <-- ASEAN ditambahkan di sini
df_yoy = self._compute_yoy_columns(df_with_asean)
df_scored = self._compute_scores(df_yoy)
df_final = self._assign_performance(df_scored)
rows_loaded = self._save(df_final)
self.pipeline_metadata["rows_loaded"] = rows_loaded
self._build_narrative_table(df_final)
self.pipeline_metadata["end_time"] = datetime.now()
duration = (
self.pipeline_metadata["end_time"] - self.pipeline_start
).total_seconds()
self.logger.info("\n" + "=" * 80)
self.logger.info("COMPLETED")
self.logger.info("=" * 80)
self.logger.info(f" Duration : {duration:.2f}s")
self.logger.info(f" Rows Fetched : {self.pipeline_metadata['rows_fetched']:,}")
self.logger.info(f" Rows Loaded (norm) : {rows_loaded:,}")
self.logger.info(f" Rows Loaded (narrative) : {self.pipeline_metadata['rows_loaded_narrative']:,}")
self.logger.info(f" sdgs_start_year : {self.sdgs_start_year}")
# =============================================================================
# AIRFLOW TASK
# =============================================================================
def run_indicator_norm_aggregation():
client = get_bigquery_client()
agg = IndicatorNormAggregator(client)
agg.run()
print(f"agg_indicator_norm loaded : {agg.pipeline_metadata['rows_loaded']:,} rows")
print(f"agg_narrative_indicator loaded: {agg.pipeline_metadata['rows_loaded_narrative']:,} rows")
# =============================================================================
# MAIN
# =============================================================================
if __name__ == "__main__":
import sys, io
if sys.stdout.encoding and sys.stdout.encoding.lower() not in ("utf-8", "utf8"):
sys.stdout = io.TextIOWrapper(sys.stdout.buffer, encoding="utf-8", errors="replace")
if sys.stderr.encoding and sys.stderr.encoding.lower() not in ("utf-8", "utf8"):
sys.stderr = io.TextIOWrapper(sys.stderr.buffer, encoding="utf-8", errors="replace")
print("=" * 80)
print("INDICATOR NORM AGGREGATION -> fs_asean_gold")
print(f" ASEAN merged into country tables (country_id={ASEAN_COUNTRY_ID})")
print("=" * 80)
logger = setup_logging()
client = get_bigquery_client()
agg = IndicatorNormAggregator(client)
agg.run()
print("\n" + "=" * 80)
print("[OK] COMPLETED")
print("=" * 80)