1 line
122 KiB
XML
1 line
122 KiB
XML
{"version":"v1alpha","kind":"VisualBuilder","metadata":{"name":"forecasting_workflow","description":"Forecasting workflow","runtime":"spark"},"spec":{"ui":{"edges":[{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"data-reader__0","target":"filter__0","id":"xy-edge__data-reader__0-filter__0"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"filter__0","target":"rest-api-invoke__0","id":"xy-edge__filter__0-rest-api-invoke__0"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"rest-api-invoke__0","target":"data-mapper__0","id":"xy-edge__rest-api-invoke__0-data-mapper__0"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"data-reader__1","target":"data-join__0","targetHandle":"target-a","id":"xy-edge__data-reader__1-data-join__0target-a"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"data-join__0","target":"filter__1","id":"xy-edge__data-join__0-filter__1"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"filter__1","target":"data-mapper__1","id":"xy-edge__filter__1-data-mapper__1"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"data-mapper__2","target":"data-writer__1","id":"xy-edge__data-mapper__2-data-writer__1"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"data-mapper__1","target":"code-transform__1","id":"xy-edge__data-mapper__1-code-transform__1"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"data-mapper__1","target":"data-writer_cloned","id":"xy-edge__data-mapper__1-data-writer_cloned"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"code-transform__1","target":"data-mapper__2","id":"xy-edge__code-transform__1-data-mapper__2"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"data-mapper__0","target":"data-mapper__3","id":"xy-edge__data-mapper__0-data-mapper__3"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"data-mapper__3","target":"data-mapper__4","id":"xy-edge__data-mapper__3-data-mapper__4"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"data-mapper__4","target":"data-mapper_cloned_2","id":"xy-edge__data-mapper__4-data-mapper_cloned_2"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"data-mapper_cloned_2","target":"data-join__0","targetHandle":"target-a","id":"xy-edge__data-mapper_cloned_2-data-join__0target-a"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"data-mapper_cloned_2","target":"code-transform__1","id":"xy-edge__data-mapper_cloned_2-code-transform__1"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"data-mapper__2","target":"data-writer_cloned","id":"xy-edge__data-mapper__2-data-writer_cloned"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"filter__1","target":"data-mapper_cloned_2_cloned","id":"xy-edge__filter__1-data-mapper_cloned_2_cloned"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"data-mapper_cloned_2_cloned","target":"code-transform__0","id":"xy-edge__data-mapper_cloned_2_cloned-code-transform__0"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"code-transform__0","target":"data-mapper_cloned_2_cloned_cloned","id":"xy-edge__code-transform__0-data-mapper_cloned_2_cloned_cloned"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"data-mapper_cloned_2_cloned_cloned","target":"code-transform_cloned","id":"xy-edge__data-mapper_cloned_2_cloned_cloned-code-transform_cloned"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"code-transform_cloned","target":"data-writer_cloned_cloned","id":"xy-edge__code-transform_cloned-data-writer_cloned_cloned"}],"nodes":[{"id":"data-reader__0","type":"workflowNode","position":{"x":28.708450583646908,"y":-362.48200237205447},"data":{"nodeType":"data-reader","id":"data-reader__0"},"measured":{"width":240,"height":122},"selected":true,"dragging":false,"style":{}},{"id":"rest-api-invoke__0","type":"workflowNode","position":{"x":708.9803552439467,"y":-359.6182330366048},"data":{"nodeType":"rest-api-invoke","id":"rest-api-invoke__0"},"measured":{"width":240,"height":122},"selected":false,"dragging":false,"style":{}},{"id":"filter__0","type":"workflowNode","position":{"x":373.78136049797183,"y":-359.7798632332295},"data":{"nodeType":"filter","id":"filter__0"},"measured":{"width":240,"height":122},"selected":false,"dragging":false,"style":{}},{"id":"data-reader__1","type":"workflowNode","position":{"x":1994.2428943628847,"y":-63.48943077545711},"data":{"nodeType":"data-reader","id":"data-reader__1"},"measured":{"width":240,"height":122},"selected":false,"dragging":false,"style":{}},{"id":"data-mapper__0","type":"workflowNode","position":{"x":1035.2257514615978,"y":-361.631330702255},"data":{"nodeType":"data-mapper","id":"data-mapper__0"},"measured":{"width":240,"height":122},"selected":false,"dragging":false,"style":{}},{"id":"data-join__0","type":"workflowNode","position":{"x":2309.8205820198277,"y":-238.1731809386307},"data":{"nodeType":"data-join","id":"data-join__0"},"measured":{"width":240,"height":122},"selected":false,"dragging":false,"style":{}},{"id":"filter__1","type":"workflowNode","position":{"x":2609.3223241261853,"y":-246.96996789894163},"data":{"nodeType":"filter","id":"filter__1"},"measured":{"width":240,"height":122},"selected":false,"dragging":false,"style":{}},{"id":"data-mapper__1","type":"workflowNode","position":{"x":2904.7217794735707,"y":-247.60998779550184},"data":{"nodeType":"data-mapper","id":"data-mapper__1"},"measured":{"width":240,"height":122},"selected":false,"dragging":false,"style":{}},{"id":"data-writer_cloned","type":"workflowNode","position":{"x":4264.760808947707,"y":-206.04133720590033},"data":{"nodeType":"data-writer","id":"data-writer_cloned"},"measured":{"width":240,"height":122},"selected":false,"dragging":false,"style":{}},{"id":"code-transform__1","type":"workflowNode","position":{"x":3362.868492902819,"y":-764.0251464836689},"data":{"nodeType":"code-transform","id":"code-transform__1"},"measured":{"width":240,"height":154},"selected":false,"dragging":false,"style":{}},{"id":"data-mapper__2","type":"workflowNode","position":{"x":3895.4642134308447,"y":-451.2543889972694},"data":{"nodeType":"data-mapper","id":"data-mapper__2"},"measured":{"width":240,"height":130},"selected":false,"dragging":false,"style":{}},{"id":"data-writer__1","type":"workflowNode","position":{"x":4278.371798251019,"y":-457.03060907637564},"data":{"nodeType":"data-writer","id":"data-writer__1"},"measured":{"width":240,"height":130},"selected":false,"dragging":false,"style":{}},{"id":"data-mapper__3","type":"workflowNode","position":{"x":1343.9319301402759,"y":-364.1540551576818},"data":{"nodeType":"data-mapper","id":"data-mapper__3"},"measured":{"width":240,"height":122},"selected":false,"dragging":false,"style":{}},{"id":"data-mapper__4","type":"workflowNode","position":{"x":1667.9656263661966,"y":-366.72377933230473},"data":{"nodeType":"data-mapper","id":"data-mapper__4"},"measured":{"width":240,"height":122},"selected":false,"dragging":false,"style":{}},{"id":"data-mapper_cloned_2","type":"workflowNode","position":{"x":1985.9753955790711,"y":-367.60998779550187},"data":{"nodeType":"data-mapper","id":"data-mapper_cloned_2"},"measured":{"width":240,"height":122},"selected":false,"dragging":false,"style":{}},{"id":"data-mapper_cloned_2_cloned","type":"workflowNode","position":{"x":2899.2214198205534,"y":112.2013106943028},"data":{"nodeType":"data-mapper","id":"data-mapper_cloned_2_cloned"},"measured":{"width":240,"height":122},"selected":false,"dragging":false,"style":{}},{"id":"code-transform__0","type":"workflowNode","position":{"x":3190.050016273345,"y":106.8695282851497},"data":{"nodeType":"code-transform","id":"code-transform__0"},"measured":{"width":240,"height":130},"selected":false,"dragging":false},{"id":"data-mapper_cloned_2_cloned_cloned","type":"workflowNode","position":{"x":3489.663363810525,"y":105.19216245384504},"data":{"nodeType":"data-mapper","id":"data-mapper_cloned_2_cloned_cloned"},"measured":{"width":240,"height":130},"selected":false,"dragging":false,"style":{}},{"id":"code-transform_cloned","type":"workflowNode","position":{"x":3768.032331064804,"y":114.53022978001482},"data":{"nodeType":"code-transform","id":"code-transform_cloned"},"measured":{"width":240,"height":122},"selected":false,"dragging":false},{"id":"data-writer_cloned_cloned","type":"workflowNode","position":{"x":4134.2630111709595,"y":90.42322868434687},"data":{"nodeType":"data-writer","id":"data-writer_cloned_cloned"},"measured":{"width":240,"height":130},"selected":false,"dragging":false,"style":{}}],"nodesData":{"data-reader__0":{"columns":[],"typeLabel":"Spark","datasetName":"","isDefault":false,"connectedComponents":[],"name":"readCustomers","type":"SparkReader","format":"iceberg","credentials":{"accessKey":"S3_ACCESS_KEY","secretKey":"S3_SECRET_KEY"},"iceberg_catalog":"dremio","region":"us-west-1","table_name":"customer"},"rest-api-invoke__0":{"name":"getBills","type":"RESTInvoke","datasource":"filterActiveCustomers","url":"https://fw-gateway:8200/fw-notification/outbound-message-config/publish","method":"POST","headers":{"Content-Type":{"value":"application/json","secret":null},"Content-type":{"value":"application/json","secret":null},"api-key":{"value":null,"secret":"OCULAR_API_KEY"},"x-tenantCode":{"value":"UTILITIES","secret":null}},"bodyTemplate":"{\n \"outMsgConfigCode\": \"EXP_ACCOUNT_BILL_HISTORY\",\n \"msgData\": {\n \"accountId\": \"{{account_id}}\",\n \"numberOfMonthPast\": \"24\"\n }\n}","materialization_strategy":"NONE","fail_on_error":true,"isDefault":false,"connectedComponents":["filterActiveCustomers"]},"filter__0":{"name":"filterActiveCustomers","type":"Filter","datasource":"readCustomers","condition":"TRIM(UPPER(status)) = \\'ACTIVE\\'","materialization_strategy":"NONE","fail_on_error":true,"isDefault":false,"connectedComponents":["readCustomers"]},"data-reader__1":{"typeLabel":"Spark","columns":[],"isDefault":false,"connectedComponents":[],"name":"readLatestBillIds","type":"SparkReader","format":"iceberg","credentials":{"accessKey":"S3_ACCESS_KEY","secretKey":"S3_SECRET_KEY"},"iceberg_catalog":"dremio","region":"us-west-1","table_name":"bills"},"data-mapper__0":{"name":"MapLatestBill","type":"DataMapping","datasource":"getBills","includeExistingColumns":false,"toSchema":[{"fieldName":"accounts","valueExpression":"from_json(\r\n get_json_object(response_body, \\'$.data\\'),\r\n \\'struct<\r\n accountId:string,\r\n numberOfMonthPast:string,\r\n output:struct<\r\n bills:array<struct<\r\n billId:string,\r\n billStatus:string,\r\n billStatusName:string,\r\n billDate:string,\r\n completionDttm:string,\r\n dueDate:string,\r\n amount:string\r\n >>\r\n >\r\n >\\'\r\n)"},{"fieldName":"account_id","valueExpression":"account_id"}],"materialization_strategy":"NONE","fail_on_error":true,"additionalData":{"isGlossaryAssisted":false,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualMappings":[{"id":"mapping-1782329233261","newFieldName":"accounts","mappingType":"expression","value":"from_json(\r\n get_json_object(response_body, '$.data'),\r\n 'struct<\r\n accountId:string,\r\n numberOfMonthPast:string,\r\n output:struct<\r\n bills:array<struct<\r\n billId:string,\r\n billStatus:string,\r\n billStatusName:string,\r\n billDate:string,\r\n completionDttm:string,\r\n dueDate:string,\r\n amount:string\r\n >>\r\n >\r\n >'\r\n)"},{"id":"mapping-1782329846646","newFieldName":"account_id","mappingType":"sourceColumn","value":"account_id"}],"manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[],"after":"","totalSourceTerms":0},"isDefault":false,"connectedComponents":["getBills"]},"data-join__0":{"name":"data_join__0","type":"RelationalJoin","dropDuplicatedColumns":true,"baseData":"data_mapper__5","joinOrder":[{"with":"readLatestBillIds","joinColumns":[{"account_Id":"account_Id"},{"mapper_bill_id":"bill_id"}],"how":"left outer"}],"isDefault":false,"connectedComponents":["readLatestBillIds","data_mapper__5"]},"filter__1":{"name":"filter__1","type":"Filter","datasource":"data_join__0","condition":"bill_id IS NULL OR mapper_bill_id <> bill_id","materialization_strategy":"NONE","fail_on_error":true,"isDefault":false,"connectedComponents":["data_join__0"]},"data-mapper__1":{"name":"BillWriterMapper","type":"DataMapping","datasource":"filter__1","includeExistingColumns":false,"toSchema":[{"fieldName":"account_id","valueExpression":"account_id"},{"fieldName":"bill_id","valueExpression":"mapper_bill_id"},{"fieldName":"bill_date","valueExpression":"bill_date"},{"fieldName":"bill_status","valueExpression":"bill_status"},{"fieldName":"due_date","valueExpression":"due_date"},{"fieldName":"created_at","valueExpression":"current_timestamp()"},{"fieldName":"id","valueExpression":"id"},{"fieldName":"amount","valueExpression":"amount"},{"fieldName":"amount_value","valueExpression":"amount_value"},{"fieldName":"bill_status_name","valueExpression":"bill_status_name"},{"fieldName":"completion_dttm","valueExpression":"completion_dttm"}],"materialization_strategy":"NONE","fail_on_error":true,"additionalData":{"isGlossaryAssisted":false,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualMappings":[{"id":"mapping-1781114585148","newFieldName":"account_id","mappingType":"sourceColumn","value":"account_id"},{"id":"mapping-1782331737547","newFieldName":"bill_id","mappingType":"sourceColumn","value":"mapper_bill_id"},{"id":"mapping-1782331770716","newFieldName":"bill_date","mappingType":"sourceColumn","value":"bill_date"},{"id":"mapping-1782331787498","newFieldName":"bill_status","mappingType":"sourceColumn","value":"bill_status"},{"id":"mapping-1782331808058","newFieldName":"due_date","mappingType":"sourceColumn","value":"due_date"},{"id":"mapping-1782331832635","newFieldName":"created_at","mappingType":"expression","value":"current_timestamp()"},{"id":"mapping-1782331863479","newFieldName":"id","mappingType":"sourceColumn","value":"id"},{"id":"mapping-1782331899605","newFieldName":"amount","mappingType":"sourceColumn","value":"amount"},{"id":"mapping-1783665970895","newFieldName":"amount_value","mappingType":"sourceColumn","value":"amount_value"},{"id":"mapping-1783665993133","newFieldName":"bill_status_name","mappingType":"sourceColumn","value":"bill_status_name"},{"id":"mapping-1784539493167","newFieldName":"completion_dttm","mappingType":"sourceColumn","value":"completion_dttm"}],"manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[],"after":"","totalSourceTerms":0},"isDefault":false,"connectedComponents":["filter__1"]},"data-writer_cloned":{"name":"data_writer__1","type":"SparkWriter","format":"iceberg","mode":"append","datasource":"BillWriterMapper","typeLabel":"Spark","credentials":{"accessKey":"S3_ACCESS_KEY","secretKey":"S3_SECRET_KEY"},"iceberg_catalog":"dremio","region":"us-west-1","table_name":"bills","isDefault":false,"connectedComponents":["BillWriterMapper","forecast_insight_data_mapper"]},"code-transform__1":{"name":"forecast_insight_code_transform","type":"CodeTransform","language":"python","datasource":"data_mapper__5","code":"try:\n\n import builtins\n import json as json_lib\n import traceback\n from datetime import datetime, timedelta\n import numpy as np\n\n # =====================================================\n # CHECK PROPHET\n # =====================================================\n\n try:\n from prophet import Prophet\n PROPHET_AVAILABLE = False\n # print(\"Prophet installed\")\n except Exception as e:\n PROPHET_AVAILABLE = True\n # print(f\"Prophet not available: {e}\")\n\n # =====================================================\n # THRESHOLDS (mirrors BillForecastService class constants)\n # =====================================================\n\n HIGH_USAGE_THRESHOLD = 15.0 # % increase -> high_usage\n DROP_THRESHOLD = -15.0 # % decrease -> drop_detected\n MIN_BILLS_FOR_FILTERING = 3 # minimum bills to apply IQR outlier filtering\n TREND_DAMPEN = 0.5 # apply only 50% of observed MoM change in fallback\n\n # =====================================================\n # READ SOURCE\n # =====================================================\n\n source_df = {{datasource}}_df\n # print(\"Input schema:\")\n source_df.printSchema()\n\n pdf = (\n source_df\n .select(\n \"account_id\",\n \"mapper_bill_id\",\n \"bill_date\",\n \"amount_value\"\n )\n .toPandas()\n )\n\n\n pdf[\"bill_date\"] = pd.to_datetime(pdf[\"bill_date\"])\n\n # print(\"Input rows =\", len(pdf))\n\n # =====================================================\n # MAPPER FILTER — only keep accounts present in mapper_df\n # =====================================================\n \n mapper_df = BillWriterMapper_df\n \n mapper_pdf = (\n mapper_df\n .select(\"account_id\")\n .toPandas()\n )\n \n if mapper_pdf.empty:\n print(\"Mapper has no data - returning empty output successfully\")\n pdf = pdf.iloc[0:0]\n else:\n mapper_account_ids = set(mapper_pdf[\"account_id\"].dropna().unique())\n print(\"Mapper account count =\", len(mapper_account_ids))\n \n before_count = len(pdf)\n pdf = pdf[pdf[\"account_id\"].isin(mapper_account_ids)].reset_index(drop=True)\n print(f\"Filtered source rows by mapper: {before_count} -> {len(pdf)}\")\n \n print(\"Input rows after mapper filter =\", len(pdf))\n\n\n output_rows = []\n\n # =====================================================\n # HELPER FUNCTIONS — outlier filtering & weighting\n # =====================================================\n\n def filter_outliers(amounts):\n \"\"\"Remove outliers using IQR method. Returns filtered list (at least 2 values kept).\"\"\"\n if len(amounts) < 3:\n return amounts\n\n sorted_vals = sorted(amounts)\n n = len(sorted_vals)\n q1 = sorted_vals[n // 4]\n q3 = sorted_vals[(3 * n) // 4]\n iqr = q3 - q1\n\n # Use 1.5x IQR rule; if IQR is 0, fall back to median +/- band\n if iqr > 0:\n lower_bound = q1 - 1.5 * iqr\n upper_bound = q3 + 1.5 * iqr\n else:\n median = sorted_vals[n // 2]\n lower_bound = median * 0.2\n upper_bound = median * 3.0\n\n filtered = [a for a in amounts if lower_bound <= a <= upper_bound]\n\n # Always keep at least the 2 most recent values\n if len(filtered) < 2:\n filtered = amounts[-2:]\n\n return filtered\n\n def exponential_weights(n, decay=0.5):\n \"\"\"Generate exponential decay weights - most recent gets highest weight.\n\n Example with n=3, decay=0.5: [0.25, 0.5, 1.0] -> normalized to [0.143, 0.286, 0.571]\n \"\"\"\n raw = [decay ** (n - 1 - i) for i in range(n)]\n total = builtins.sum(raw)\n return [w / total for w in raw]\n\n def detect_anomalies(amounts, threshold=2.0):\n \"\"\"Detect anomalies using Z-score method.\"\"\"\n if len(amounts) < 3:\n return [False] * len(amounts)\n\n mean_val = np.mean(amounts)\n std_val = np.std(amounts)\n\n if std_val == 0:\n return [False] * len(amounts)\n\n z_scores = [(x - mean_val) / std_val for x in amounts]\n return [bool(abs(z) > threshold) for z in z_scores]\n\n def calculate_trend_slope(amounts):\n \"\"\"Calculate normalized trend slope using linear regression (% change per period).\"\"\"\n if len(amounts) < 2:\n return 0.0\n\n x = np.arange(len(amounts))\n y = np.array(amounts)\n\n n = len(x)\n denom = (n * np.sum(x ** 2) - np.sum(x) ** 2)\n if denom == 0:\n return 0.0\n slope = (n * np.sum(x * y) - np.sum(x) * np.sum(y)) / denom\n\n mean_val = np.mean(amounts)\n if mean_val > 0:\n return (slope / mean_val) * 100\n return 0.0\n\n # =====================================================\n # HELPER FUNCTIONS — forecasting\n # =====================================================\n\n def prophet_forecast(prophet_df, periods=3):\n \"\"\"Run Prophet forecast for the next `periods` months. Raises on failure\n so the caller can fall back to fallback_forecast_weighted.\"\"\"\n model = Prophet(\n yearly_seasonality=True,\n weekly_seasonality=False,\n daily_seasonality=False,\n interval_width=0.80\n )\n\n model.fit(prophet_df)\n\n future = model.make_future_dataframe(\n periods=periods,\n freq=\"M\"\n )\n\n pred = model.predict(future)\n\n forecast_df = pred[\n pred[\"ds\"] > prophet_df[\"ds\"].max()\n ][[\n \"ds\",\n \"yhat\",\n \"yhat_lower\",\n \"yhat_upper\"\n ]].copy()\n\n forecast_df[\"yhat\"] = forecast_df[\"yhat\"].clip(lower=0).round(2)\n forecast_df[\"yhat_lower\"] = forecast_df[\"yhat_lower\"].clip(lower=0).round(2)\n forecast_df[\"yhat_upper\"] = forecast_df[\"yhat_upper\"].round(2)\n\n return forecast_df\n\n def fallback_forecast_weighted(prophet_df, periods=3):\n \"\"\"\n Weighted-average fallback with dampened trend (used when Prophet is\n unavailable, fails, or there are fewer than 4 data points).\n\n Steps:\n 1. Outlier filtering (IQR method) when >= MIN_BILLS_FOR_FILTERING bills.\n 2. Exponential decay weighting (most recent bill weighted highest).\n 3. Dampened month-over-month trend projection (50% of observed rate),\n with a seasonal override when a same-calendar-month average exists.\n\n Confidence interval: +/-15% around the forecast value.\n \"\"\"\n amounts_all = [float(v) for v in prophet_df[\"y\"].values]\n dates_all = list(prophet_df[\"ds\"].values)\n\n amounts = [a for a in amounts_all if a > 0]\n if not amounts:\n return pd.DataFrame(columns=[\"ds\", \"yhat\", \"yhat_lower\", \"yhat_upper\"])\n\n # Step 1: outlier filtering\n if len(amounts) >= MIN_BILLS_FOR_FILTERING:\n clean_amounts = filter_outliers(amounts)\n else:\n clean_amounts = amounts\n\n # Seasonal map: month-of-year -> list of historical amounts in that month\n monthly_map = {}\n for d, a in zip(dates_all, amounts_all):\n if a <= 0:\n continue\n month = pd.Timestamp(d).month\n monthly_map.setdefault(month, []).append(a)\n\n last_date = prophet_df[\"ds\"].max()\n\n # Step 2: exponential decay weighted average on clean data\n weights = exponential_weights(len(clean_amounts))\n weighted_avg = builtins.sum(a * w for a, w in zip(clean_amounts, weights))\n\n # Step 3: dampened month-over-month trend\n if len(clean_amounts) >= 2:\n mom_changes = []\n for j in range(1, len(clean_amounts)):\n if clean_amounts[j - 1] > 0:\n mom_changes.append(\n (clean_amounts[j] - clean_amounts[j - 1]) / clean_amounts[j - 1]\n )\n avg_mom = (builtins.sum(mom_changes) / len(mom_changes)) if mom_changes else 0.0\n dampened_mom = avg_mom * TREND_DAMPEN\n else:\n dampened_mom = 0.0\n\n rows = []\n base_val = weighted_avg\n\n for i in range(1, periods + 1):\n future_dt = last_date + pd.DateOffset(months=i)\n future_month = future_dt.month\n\n if future_month in monthly_map and monthly_map[future_month]:\n seasonal_avg = builtins.sum(monthly_map[future_month]) / len(monthly_map[future_month])\n predicted_value = seasonal_avg\n else:\n predicted_value = builtins.max(0.0, base_val * (1 + dampened_mom) ** i)\n\n lower_bound = builtins.max(0.0, predicted_value * 0.85)\n upper_bound = predicted_value * 1.15\n\n rows.append({\n \"ds\": future_dt,\n \"yhat\": round(predicted_value, 2),\n \"yhat_lower\": round(lower_bound, 2),\n \"yhat_upper\": round(upper_bound, 2)\n })\n\n # print(\n # f\"Fallback forecast: {len(amounts)} bills -> {len(clean_amounts)} clean -> \"\n # f\"base ${weighted_avg:.2f}, dampened MoM {dampened_mom * 100:.1f}%\"\n # )\n\n return pd.DataFrame(rows)\n\n # =====================================================\n # HELPER FUNCTIONS — classification, severity, explanation\n # =====================================================\n\n def classify_type(recent_amounts, forecast_amounts):\n \"\"\"Classify insight type based on % change between recent avg and forecast avg.\"\"\"\n if not recent_amounts or not forecast_amounts:\n return \"stable_usage\"\n\n recent_avg = builtins.sum(recent_amounts) / len(recent_amounts)\n forecast_avg = builtins.sum(forecast_amounts) / len(forecast_amounts)\n\n if recent_avg == 0:\n return \"stable_usage\"\n\n pct_change = ((forecast_avg - recent_avg) / recent_avg) * 100\n\n if pct_change >= HIGH_USAGE_THRESHOLD:\n return \"high_usage\"\n elif pct_change <= DROP_THRESHOLD:\n return \"drop_detected\"\n else:\n return \"stable_usage\"\n\n def compute_severity_score(recent_amounts, forecast_amounts):\n \"\"\"Compute a 1-10 severity score based on magnitude of change.\"\"\"\n if not recent_amounts or not forecast_amounts:\n return 1\n\n recent_avg = builtins.sum(recent_amounts) / len(recent_amounts)\n forecast_avg = builtins.sum(forecast_amounts) / len(forecast_amounts)\n\n if recent_avg == 0:\n return 1\n\n pct_change = abs(((forecast_avg - recent_avg) / recent_avg) * 100)\n return builtins.min(10, builtins.max(1, int(pct_change / 10) + 1))\n\n def generate_explanation_template(recent_amounts, forecast_amounts, insight_type, bill_count):\n \"\"\"Template-based alert/message/explanation for the Rank 1 forecast insight.\"\"\"\n recent_avg = builtins.sum(recent_amounts) / len(recent_amounts) if recent_amounts else 0\n forecast_avg = builtins.sum(forecast_amounts) / len(forecast_amounts) if forecast_amounts else 0\n\n if recent_avg > 0:\n pct_change = ((forecast_avg - recent_avg) / recent_avg) * 100\n else:\n pct_change = 0\n\n direction = \"increase\" if pct_change > 0 else \"decrease\"\n abs_pct = abs(pct_change)\n\n type_labels = {\n \"high_usage\": \"High Usage Expected\",\n \"drop_detected\": \"Bill Drop Detected\",\n \"stable_usage\": \"Stable Billing Pattern\",\n }\n\n alert = type_labels.get(insight_type, \"Bill Forecast\")\n message = f\"Forecasted bills show a {abs_pct:.1f}% {direction} over the next 3 months.\"\n explanation = (\n f\"Based on the last {bill_count} months of billing data, \"\n f\"the average recent bill is ${recent_avg:.2f} and the forecasted average is ${forecast_avg:.2f}. \"\n f\"This represents a {abs_pct:.1f}% {direction} \"\n f\"(${abs(forecast_avg - recent_avg):.2f} difference).\"\n )\n\n return alert, message, explanation\n\n # =====================================================\n # HELPER FUNCTIONS — graphs & considered bills\n # =====================================================\n\n def build_bar_graph(labels, data, label, color=\"rgba(75,192,192,0.6)\"):\n return {\n \"type\": \"bar\",\n \"labels\": labels,\n \"datasets\": [{\n \"label\": label,\n \"data\": data,\n \"backgroundColor\": color\n }]\n }\n\n def build_forecast_graph(forecast_labels, forecast_values, forecast_lower, forecast_upper):\n return {\n \"type\": \"bar\",\n \"labels\": forecast_labels,\n \"datasets\": [\n {\n \"label\": \"Forecasted Bills\",\n \"data\": forecast_values,\n \"borderColor\": \"rgb(75,102,192)\",\n \"backgroundColor\": \"rgba(75,192,192,0.2)\",\n \"fill\": True\n },\n {\n \"label\": \"Confidence Lower\",\n \"data\": forecast_lower,\n \"borderColor\": \"rgba(75,192,192,0.3)\",\n \"backgroundColor\": \"transparent\",\n \"borderDash\": [5, 5],\n \"fill\": False\n },\n {\n \"label\": \"Confidence Upper\",\n \"data\": forecast_upper,\n \"borderColor\": \"rgba(75,192,192,0.3)\",\n \"backgroundColor\": \"transparent\",\n \"borderDash\": [5, 5],\n \"fill\": False\n }\n ]\n }\n\n def build_considered_bills(rows_df, anomaly_flags=None):\n \"\"\"Build the consideredBills list from a pandas slice of bill rows.\"\"\"\n considered = []\n for i, (_, r) in enumerate(rows_df.iterrows()):\n is_anomaly = bool(anomaly_flags[i]) if anomaly_flags and i < len(anomaly_flags) else False\n considered.append({\n \"billID\": str(r[\"mapper_bill_id\"]),\n \"billDate\": r[\"bill_date\"].strftime(\"%Y-%m-%d\"),\n \"billAmount\": round(float(r[\"amount_value\"]), 2),\n \"consumptionValue\": round(float(r[\"amount_value\"]), 2),\n \"consumptionUnit\": \"USD\",\n \"isAnomaly\": is_anomaly\n })\n return considered\n\n # =====================================================\n # INSIGHT BUILDER — Rank 1: Bill Forecast\n # =====================================================\n\n def build_forecast_insight(recent, forecast_df, anomaly_flags):\n \"\"\"\n Rank 1 insight: forecast classification (high_usage / drop_detected /\n stable_usage), with historical graph + forecast graph + consideredBills.\n \"\"\"\n actual_labels = [d.strftime(\"%Y-%m\") for d in recent[\"bill_date\"]]\n actual_amounts = [round(float(x), 2) for x in recent[\"amount_value\"]]\n\n forecast_labels = [d.strftime(\"%Y-%m\") for d in forecast_df[\"ds\"]]\n forecast_values = [round(float(x), 2) for x in forecast_df[\"yhat\"]]\n forecast_lower = [round(float(x), 2) for x in forecast_df[\"yhat_lower\"]]\n forecast_upper = [round(float(x), 2) for x in forecast_df[\"yhat_upper\"]]\n\n insight_type = classify_type(actual_amounts, forecast_values)\n severity = compute_severity_score(actual_amounts, forecast_values)\n alert, message, explanation = generate_explanation_template(\n actual_amounts, forecast_values, insight_type, len(recent)\n )\n\n considered_bills = build_considered_bills(recent, anomaly_flags)\n\n actual_graph = build_bar_graph(actual_labels, actual_amounts, \"Bills USD\")\n forecast_graph = build_forecast_graph(forecast_labels, forecast_values, forecast_lower, forecast_upper)\n\n insight = {\n \"rank\": 1,\n \"alert\": alert,\n \"message\": message,\n \"explanation\": explanation,\n \"severityScore\": severity,\n \"consideredBills\": considered_bills,\n \"graph\": actual_graph,\n \"forecastGraph\": forecast_graph,\n \"type\": insight_type\n }\n\n return insight, forecast_labels, forecast_values, forecast_lower, forecast_upper\n\n # =====================================================\n # INSIGHT BUILDER — Rank 2: Trend Summary (last 3 months)\n # =====================================================\n\n def build_trend_insight(recent, anomaly_flags, forecast_labels, forecast_values, forecast_lower, forecast_upper):\n \"\"\"\n Rank 2 insight: month-over-month trend pattern across the last 3 bills.\n\n Patterns:\n declining_trend - both MoM changes < -20%\n increasing_trend - both MoM changes > +20%\n spike_resolved - oldest month 30%+ higher, bills dropped since\n mid_spike - middle month 30%+ higher than neighbors\n recent_spike - most recent month jumped 30%+\n stable_trend - all within 20% of 3-month average\n\n Returns None if fewer than 3 bills or no clear pattern.\n \"\"\"\n if len(recent) < 3:\n return None\n\n amounts = [round(float(x), 2) for x in recent[\"amount_value\"]]\n dates = [d.strftime(\"%Y-%m\") for d in recent[\"bill_date\"]]\n month_names = [d.strftime(\"%B %Y\") for d in recent[\"bill_date\"]]\n\n a0, a1, a2 = amounts # oldest -> newest\n\n def pct(old, new):\n return ((new - old) / old * 100) if old != 0 else 0\n\n chg_1 = pct(a0, a1)\n chg_2 = pct(a1, a2)\n total_chg = pct(a0, a2)\n\n peak_idx = amounts.index(builtins.max(amounts))\n\n alert = \"\"\n message = \"\"\n explanation = \"\"\n trend_type = \"stable_trend\"\n\n if chg_1 < -20 and chg_2 < -20:\n trend_type = \"declining_trend\"\n alert = \"Bills Declining Steadily\"\n message = (\n f\"Your bill has dropped {abs(total_chg):.0f}% over the last 3 months \"\n f\"- from ${a0:,.2f} in {month_names[0]} to ${a2:,.2f} in {month_names[2]}.\"\n )\n explanation = (\n f\"{month_names[0]}: ${a0:,.2f} -> {month_names[1]}: ${a1:,.2f} ({chg_1:+.1f}%) -> \"\n f\"{month_names[2]}: ${a2:,.2f} ({chg_2:+.1f}%). \"\n f\"This is a consistent downward trend that may continue.\"\n )\n\n elif chg_1 > 20 and chg_2 > 20:\n trend_type = \"increasing_trend\"\n alert = \"Bills Increasing Steadily\"\n message = (\n f\"Your bill has risen {abs(total_chg):.0f}% over the last 3 months \"\n f\"- from ${a0:,.2f} in {month_names[0]} to ${a2:,.2f} in {month_names[2]}.\"\n )\n explanation = (\n f\"{month_names[0]}: ${a0:,.2f} -> {month_names[1]}: ${a1:,.2f} ({chg_1:+.1f}%) -> \"\n f\"{month_names[2]}: ${a2:,.2f} ({chg_2:+.1f}%). \"\n f\"Your usage has been climbing; consider reviewing recent activity.\"\n )\n\n elif peak_idx == 0 and abs(total_chg) > 30:\n trend_type = \"spike_resolved\"\n alert = \"Recent Bill Spike Has Resolved\"\n message = (\n f\"Your bill was ${a0:,.2f} in {month_names[0]} but has since dropped \"\n f\"to ${a2:,.2f} in {month_names[2]}.\"\n )\n explanation = (\n f\"The {month_names[0]} bill (${a0:,.2f}) was significantly higher than the recent \"\n f\"{month_names[1]} (${a1:,.2f}) and {month_names[2]} (${a2:,.2f}). \"\n f\"This suggests the spike was a one-time event and bills are normalizing.\"\n )\n\n elif peak_idx == 1 and pct(a1, a0) < -30 and pct(a1, a2) < -30:\n trend_type = \"mid_spike\"\n alert = f\"Bill Spike in {month_names[1]}\"\n message = (\n f\"Your {month_names[1]} bill spiked to ${a1:,.2f} but has returned \"\n f\"to ${a2:,.2f} in {month_names[2]}.\"\n )\n explanation = (\n f\"{month_names[0]}: ${a0:,.2f} -> {month_names[1]}: ${a1:,.2f} \"\n f\"(spike of {pct(a0, a1):+.1f}%) -> \"\n f\"{month_names[2]}: ${a2:,.2f} (back to {pct(a1, a2):+.1f}%). \"\n f\"The {month_names[1]} spike appears to be an anomaly.\"\n )\n\n elif peak_idx == 2 and pct(a1, a2) > 30:\n trend_type = \"recent_spike\"\n alert = \"Recent Bill Spike\"\n message = (\n f\"Your latest bill in {month_names[2]} jumped to ${a2:,.2f} \"\n f\"- up {pct(a1, a2):.0f}% from {month_names[1]} (${a1:,.2f}).\"\n )\n explanation = (\n f\"{month_names[0]}: ${a0:,.2f} -> {month_names[1]}: ${a1:,.2f} -> \"\n f\"{month_names[2]}: ${a2:,.2f} ({pct(a1, a2):+.1f}%). \"\n f\"This recent increase is worth monitoring.\"\n )\n\n else:\n avg_3 = builtins.sum(amounts) / 3\n max_dev = builtins.max(abs(a - avg_3) / avg_3 * 100 for a in amounts) if avg_3 > 0 else 0\n if max_dev < 20:\n trend_type = \"stable_trend\"\n alert = \"Bills Are Stable\"\n message = f\"Your bills have been consistent over the last 3 months, averaging ${avg_3:,.2f}.\"\n explanation = (\n f\"{month_names[0]}: ${a0:,.2f}, {month_names[1]}: ${a1:,.2f}, {month_names[2]}: ${a2:,.2f}. \"\n f\"Variation is within normal range.\"\n )\n else:\n return None\n\n severity = builtins.min(10, builtins.max(1, int(abs(total_chg) / 15) + 1))\n\n considered_bills = build_considered_bills(recent, anomaly_flags)\n\n trend_graph = {\n \"type\": \"bar\",\n \"labels\": dates,\n \"datasets\": [{\n \"label\": \"Monthly Bills USD\",\n \"data\": amounts,\n \"backgroundColor\": [\n \"rgba(255,99,132,0.6)\" if i == peak_idx else \"rgba(75,192,192,0.6)\"\n for i in range(3)\n ]\n }]\n }\n\n forecast_graph = build_forecast_graph(forecast_labels, forecast_values, forecast_lower, forecast_upper)\n\n return {\n \"rank\": 2,\n \"alert\": alert,\n \"message\": message,\n \"explanation\": explanation,\n \"severityScore\": severity,\n \"consideredBills\": considered_bills,\n \"graph\": trend_graph,\n \"forecastGraph\": forecast_graph,\n \"type\": trend_type\n }\n\n # =====================================================\n # INSIGHT BUILDER — Rank 3: Year-over-Year Comparison\n # =====================================================\n\n def build_yoy_insight(acct_df, anomaly_flags_full):\n \"\"\"\n Rank 3 insight: compares the most recent 3 months against the same 3\n calendar months a year ago. Requires >= 6 months of history overall,\n and requires that the same-month-last-year data actually exists.\n Returns None if insufficient data.\n \"\"\"\n if len(acct_df) < 6:\n return None\n\n sorted_df = acct_df.sort_values(\"bill_date\").reset_index(drop=True)\n recent_3 = sorted_df.tail(3)\n recent_dates = [d for d in recent_3[\"bill_date\"]]\n\n yoy_targets = [(d.year - 1, d.month) for d in recent_dates]\n\n yoy_rows = sorted_df[\n sorted_df[\"bill_date\"].apply(lambda d: (d.year, d.month) in yoy_targets)\n ]\n\n if len(yoy_rows) < len(yoy_targets):\n return None\n\n yoy_rows = yoy_rows.sort_values(\"bill_date\").tail(len(yoy_targets))\n\n recent_amounts = [round(float(x), 2) for x in recent_3[\"amount_value\"]]\n yoy_amounts = [round(float(x), 2) for x in yoy_rows[\"amount_value\"]]\n\n recent_avg = builtins.sum(recent_amounts) / len(recent_amounts)\n yoy_avg = builtins.sum(yoy_amounts) / len(yoy_amounts)\n pct_change = ((recent_avg - yoy_avg) / yoy_avg * 100) if yoy_avg != 0 else 0\n\n direction = \"increased\" if pct_change > 0 else \"decreased\"\n if pct_change > HIGH_USAGE_THRESHOLD:\n insight_type = \"high_usage\"\n elif pct_change < DROP_THRESHOLD:\n insight_type = \"drop_detected\"\n else:\n insight_type = \"stable_usage\"\n\n severity = builtins.min(10, builtins.max(1, int(abs(pct_change) / 10) + 1))\n\n # anomaly flags computed over the full account history align by position\n recent_anomalies = anomaly_flags_full[-3:] if len(anomaly_flags_full) >= 3 else [False] * 3\n considered_bills = build_considered_bills(recent_3, recent_anomalies) + build_considered_bills(yoy_rows, None)\n\n yoy_graph = build_bar_graph(\n [d.strftime(\"%Y-%m\") for d in yoy_rows[\"bill_date\"]],\n yoy_amounts,\n \"Last Year Bills USD\",\n color=\"rgba(153,102,255,0.6)\"\n )\n current_graph = build_bar_graph(\n [d.strftime(\"%Y-%m\") for d in recent_3[\"bill_date\"]],\n recent_amounts,\n \"Current Year Bills USD\",\n color=\"rgba(75,192,192,0.6)\"\n )\n\n return {\n \"rank\": 3,\n \"alert\": f\"Year-over-Year Bill {direction.capitalize()}\",\n \"message\": f\"Your bills have {direction} by {abs(pct_change):.1f}% compared to the same period last year.\",\n \"explanation\": (\n f\"Average bill for the recent 3 months: ${recent_avg:.2f}. \"\n f\"Average bill for the same 3 months last year: ${yoy_avg:.2f}. \"\n f\"That is a {abs(pct_change):.1f}% {direction}.\"\n ),\n \"severityScore\": severity,\n \"consideredBills\": considered_bills,\n \"graph\": yoy_graph,\n \"forecastGraph\": current_graph,\n \"type\": insight_type\n }\n\n # =====================================================\n # FORECAST ACCURACY — best-effort in-sample backtest\n # =====================================================\n\n def compute_forecast_accuracy_backtest(acct_df):\n \"\"\"\n Best-effort forecast accuracy, computed entirely from the bills already\n present in this dataframe (no external previous-forecast input available).\n\n Approach: hold out the most recent actual bill, forecast 1 month ahead\n using only the months before it (same Prophet/fallback logic as the\n live forecast), then compare that 1-month-ahead prediction against the\n real bill that came in. This approximates \"how accurate was last\n month's forecast\" without needing a stored previous forecast.\n\n Requires >= 4 bills (3 to forecast from + 1 actual to validate against).\n Returns None if not enough data.\n \"\"\"\n sorted_df = acct_df.sort_values(\"bill_date\").reset_index(drop=True)\n if len(sorted_df) < 4:\n return None\n\n train_df = sorted_df.iloc[:-1]\n actual_row = sorted_df.iloc[-1]\n\n train_prophet_df = train_df.rename(columns={\"bill_date\": \"ds\", \"amount_value\": \"y\"})[[\"ds\", \"y\"]]\n\n try:\n if PROPHET_AVAILABLE and len(train_df) >= 4:\n bt_forecast_df = prophet_forecast(train_prophet_df, periods=1)\n else:\n raise Exception(\"Prophet unavailable or insufficient data for backtest\")\n except Exception:\n bt_forecast_df = fallback_forecast_weighted(train_prophet_df, periods=1)\n\n if bt_forecast_df.empty:\n return None\n\n predicted = float(bt_forecast_df.iloc[0][\"yhat\"])\n actual = float(actual_row[\"amount_value\"])\n\n error_pct = abs(predicted - actual) / actual * 100 if actual != 0 else 0\n accuracy = builtins.max(0, 100 - error_pct)\n\n return {\n \"method\": \"in_sample_backtest\",\n \"validatedMonth\": actual_row[\"bill_date\"].strftime(\"%Y-%m\"),\n \"predicted\": round(predicted, 2),\n \"actual\": round(actual, 2),\n \"accuracyPct\": round(accuracy, 1)\n }\n\n # =====================================================\n # PROCESS EACH ACCOUNT\n # =====================================================\n\n for account_id, acct_df in pdf.groupby(\"account_id\"):\n\n try:\n\n # print(f\"\\nProcessing account {account_id}\")\n\n acct_df = acct_df.sort_values(\"bill_date\").reset_index(drop=True)\n\n bill_count = len(acct_df)\n\n # print(\"Bill count =\", bill_count)\n\n if bill_count < 3:\n # print(\"Skipping account - less than 3 bills\")\n continue\n\n # =============================================\n # PROPHET INPUT\n # =============================================\n\n prophet_df = acct_df.rename(\n columns={\n \"bill_date\": \"ds\",\n \"amount_value\": \"y\"\n }\n )[[\"ds\", \"y\"]]\n\n # =============================================\n # FORECAST (Prophet >= 4 points, else weighted fallback)\n # =============================================\n\n try:\n\n if PROPHET_AVAILABLE and bill_count >= 4:\n\n # print(\"Running Prophet\")\n forecast_df = prophet_forecast(prophet_df, periods=3)\n\n else:\n raise Exception(\"Prophet unavailable or insufficient data (<4 points)\")\n\n except Exception as prophet_error:\n\n # print(f\"Prophet failed for {account_id}: {prophet_error}\")\n # print(\"Using weighted-average fallback (outlier filtering + exponential decay + dampened trend)\")\n\n forecast_df = fallback_forecast_weighted(prophet_df, periods=3)\n\n # print(\"Forecast rows =\", len(forecast_df))\n # print(\"forecast_df\", forecast_df)\n\n # =============================================\n # ANOMALY DETECTION (full history, Z-score)\n # =============================================\n\n all_amounts = [round(float(x), 2) for x in acct_df[\"amount_value\"]]\n anomaly_flags_full = detect_anomalies(all_amounts)\n\n # =============================================\n # ACTUAL DATA — recent 3 months\n # =============================================\n\n recent = acct_df.tail(3).reset_index(drop=True)\n recent_anomaly_flags = anomaly_flags_full[-3:] if len(anomaly_flags_full) >= 3 else [False] * len(recent)\n\n # print(\"recent\", recent)\n\n # =============================================\n # TREND SLOPE (informational, kept in output)\n # =============================================\n\n recent_amounts_for_slope = [round(float(x), 2) for x in recent[\"amount_value\"]]\n trend_slope = calculate_trend_slope(recent_amounts_for_slope)\n\n # =============================================\n # RANK 1 — FORECAST INSIGHT\n # =============================================\n\n forecast_insight, forecast_labels, forecast_values, forecast_lower, forecast_upper = (\n build_forecast_insight(recent, forecast_df, recent_anomaly_flags)\n )\n\n insights = [forecast_insight]\n\n # =============================================\n # RANK 2 — TREND SUMMARY\n # =============================================\n\n trend_insight = build_trend_insight(\n recent, recent_anomaly_flags,\n forecast_labels, forecast_values, forecast_lower, forecast_upper\n )\n if trend_insight:\n insights.append(trend_insight)\n\n # =============================================\n # RANK 3 — YEAR-OVER-YEAR COMPARISON\n # =============================================\n\n yoy_insight = build_yoy_insight(acct_df, anomaly_flags_full)\n if yoy_insight:\n insights.append(yoy_insight)\n\n # =============================================\n # FORECAST ACCURACY — best-effort backtest\n # =============================================\n\n accuracy_data = compute_forecast_accuracy_backtest(acct_df)\n if accuracy_data:\n for ins in insights:\n if ins.get(\"rank\") == 1:\n ins[\"previousForecastAccuracy\"] = accuracy_data\n\n # =============================================\n # FINAL JSON\n # =============================================\n\n response = {\n \"accountId\": account_id,\n \"generatedAt\": datetime.utcnow().isoformat(),\n \"cacheHit\": False,\n \"dataPointsUsed\": len(recent),\n \"nextRefreshDate\": (\n datetime.utcnow() + pd.DateOffset(months=1)\n ).strftime(\"%Y-%m-%d\"),\n \"forecastMonth\": datetime.utcnow().strftime(\"%Y-%m\"),\n \"stale\": False,\n \"forecastMethod\": \"prophet\" if (PROPHET_AVAILABLE and bill_count >= 4) else \"weighted_average_fallback\",\n \"trendSlope\": round(trend_slope, 2),\n \"insights\": insights\n }\n\n output_rows.append(\n Row(\n account_id=str(account_id),\n insight_json=json_lib.dumps(response)\n )\n )\n\n except Exception as account_error:\n\n # print(f\"Account failed {account_id}: {account_error}\")\n\n traceback.print_exc()\n continue\n\n # =====================================================\n # OUTPUT\n # =====================================================\n\n if len(output_rows) > 0:\n\n forecast_insight_code_transform_df = (\n spark.createDataFrame(output_rows)\n )\n\n else:\n\n empty_schema = (\n \"account_id string,\"\n \" insight_json string\"\n )\n\n forecast_insight_code_transform_df = (\n spark.createDataFrame(\n [],\n empty_schema\n )\n )\n\n # print(f\"Generated insights for {len(output_rows)} accounts\")\n\n # forecast_insight_code_transform_df.show(truncate=False)\n\n forecast_insight_code_transform_df.createOrReplaceTempView(\"{{name}}_df\")\n\n forecast_insight_code_transform_execute_status = \"SUCCESS\"\n\nexcept Exception as e:\n\n print(\"Pipeline failed\")\n print(str(e))\n\n forecast_insight_code_transform_execute_status = \"ERROR\"\n\n raise","materialization_strategy":"NONE","fail_on_error":true,"isDefault":false,"connectedComponents":["BillWriterMapper","data_mapper__5"]},"data-mapper__2":{"name":"forecast_insight_data_mapper","type":"DataMapping","datasource":"forecast_insight_code_transform","includeExistingColumns":false,"toSchema":[{"fieldName":"account_id","valueExpression":"account_id"},{"fieldName":"id","valueExpression":"uuid()"},{"fieldName":"insight_data","valueExpression":"get_json_object(insight_json, \\'$.insights\\')"},{"fieldName":"forecast_month","valueExpression":"get_json_object(insight_json, \\'$.forecastMonth\\')"},{"fieldName":"type","valueExpression":"get_json_object(insight_json, \\'$.insights[0].type\\')"},{"fieldName":"created_at","valueExpression":"date_format(current_timestamp(), \"yyyy-MM-dd\\'T\\'HH:mm:ss\")"},{"fieldName":"updated_at","valueExpression":"date_format(current_timestamp(), \"yyyy-MM-dd\\'T\\'HH:mm:ss\")"}],"materialization_strategy":"NONE","fail_on_error":true,"additionalData":{"isGlossaryAssisted":false,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualMappings":[{"id":"mapping-1782062869403","newFieldName":"account_id","mappingType":"sourceColumn","value":"account_id"},{"id":"mapping-1782155700948","newFieldName":"id","mappingType":"expression","value":"uuid()"},{"id":"mapping-1782156149735","newFieldName":"insight_data","mappingType":"expression","value":"get_json_object(insight_json, '$.insights')"},{"id":"mapping-1782156157631","newFieldName":"forecast_month","mappingType":"expression","value":"get_json_object(insight_json, '$.forecastMonth')"},{"id":"mapping-1782156764648","newFieldName":"type","mappingType":"expression","value":"get_json_object(insight_json, '$.insights[0].type')"},{"id":"mapping-1783671861263","newFieldName":"created_at","mappingType":"expression","value":"date_format(current_timestamp(), \"yyyy-MM-dd'T'HH:mm:ss\")"},{"id":"mapping-1783671876510","newFieldName":"updated_at","mappingType":"expression","value":"date_format(current_timestamp(), \"yyyy-MM-dd'T'HH:mm:ss\")"}],"manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[],"after":"","totalSourceTerms":0},"isDefault":false,"connectedComponents":["forecast_insight_code_transform"]},"data-writer__1":{"name":"accountinsights_data_writer","type":"SparkWriter","format":"iceberg","mode":"append","datasource":"forecast_insight_data_mapper","typeLabel":"Spark","credentials":{"accessKey":"S3_ACCESS_KEY","secretKey":"S3_SECRET_KEY"},"iceberg_catalog":"dremio","region":"us-west-1","table_name":"insights","isDefault":false,"connectedComponents":["forecast_insight_data_mapper"]},"data-mapper__3":{"name":"data_mapper__3","type":"DataMapping","datasource":"MapLatestBill","includeExistingColumns":false,"toSchema":[{"fieldName":"account_id","valueExpression":"account_id"},{"fieldName":"bills","valueExpression":"accounts.output.bills"}],"materialization_strategy":"NONE","fail_on_error":true,"additionalData":{"isGlossaryAssisted":false,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualMappings":[{"id":"mapping-1782329935585","newFieldName":"account_id","mappingType":"sourceColumn","value":"account_id"},{"id":"mapping-1782330182362","newFieldName":"bills","mappingType":"expression","value":"accounts.output.bills"}],"manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[],"after":"","totalSourceTerms":0},"isDefault":false,"connectedComponents":["MapLatestBill"]},"data-mapper__4":{"name":"data_mapper__4","type":"DataMapping","datasource":"data_mapper__3","includeExistingColumns":false,"toSchema":[{"fieldName":"account_id","valueExpression":"account_id"},{"fieldName":"final_bills","valueExpression":"explode(bills)"}],"materialization_strategy":"NONE","fail_on_error":true,"additionalData":{"isGlossaryAssisted":false,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualMappings":[{"id":"mapping-1782330221700","newFieldName":"account_id","mappingType":"sourceColumn","value":"account_id"},{"id":"mapping-1782330259011","newFieldName":"final_bills","mappingType":"expression","value":"explode(bills)"}],"manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[],"after":"","totalSourceTerms":0},"isDefault":false,"connectedComponents":["data_mapper__3"]},"data-mapper_cloned_2":{"name":"data_mapper__5","type":"DataMapping","datasource":"data_mapper__4","includeExistingColumns":false,"toSchema":[{"fieldName":"account_id","valueExpression":"account_id"},{"fieldName":"created_at","valueExpression":"current_timestamp()"},{"fieldName":"id","valueExpression":"uuid()"},{"fieldName":"bill_date","valueExpression":"to_date(final_bills.billDate)"},{"fieldName":"due_date","valueExpression":"to_date(final_bills.dueDate)"},{"fieldName":"bill_status","valueExpression":"final_bills.billStatus"},{"fieldName":"mapper_bill_id","valueExpression":"final_bills.billId"},{"fieldName":"amount_value","valueExpression":"cast(replace(replace(final_bills.amount, \\'$\\', \\'\\'), \\',\\', \\'\\')AS decimal(10, 2))"},{"fieldName":"amount","valueExpression":"final_bills.amount"},{"fieldName":"bill_status_name","valueExpression":"final_bills.billStatusName"},{"fieldName":"completion_dttm","valueExpression":"final_bills.completionDttm"}],"materialization_strategy":"NONE","fail_on_error":true,"additionalData":{"isGlossaryAssisted":false,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualMappings":[{"id":"mapping-1781114585148","newFieldName":"account_id","mappingType":"sourceColumn","value":"account_id"},{"id":"mapping-1781263224943","newFieldName":"created_at","mappingType":"expression","value":"current_timestamp()"},{"id":"mapping-1782190167525","newFieldName":"id","mappingType":"expression","value":"uuid()"},{"id":"mapping-1782331048880","newFieldName":"bill_date","mappingType":"expression","value":"to_date(final_bills.billDate)"},{"id":"mapping-1782331060005","newFieldName":"due_date","mappingType":"expression","value":"to_date(final_bills.dueDate)"},{"id":"mapping-1782331069120","newFieldName":"bill_status","mappingType":"expression","value":"final_bills.billStatus"},{"id":"mapping-1782331518221","newFieldName":"mapper_bill_id","mappingType":"expression","value":"final_bills.billId"},{"id":"mapping-1783664487911","newFieldName":"amount_value","mappingType":"expression","value":"cast(replace(replace(final_bills.amount, '$', ''), ',', '')AS decimal(10, 2))"},{"id":"mapping-1783664496058","newFieldName":"amount","mappingType":"expression","value":"final_bills.amount"},{"id":"mapping-1783664524438","newFieldName":"bill_status_name","mappingType":"expression","value":"final_bills.billStatusName"},{"id":"mapping-1783664548606","newFieldName":"completion_dttm","mappingType":"expression","value":"final_bills.completionDttm"}],"manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[],"after":"","totalSourceTerms":0},"isDefault":false,"connectedComponents":["data_mapper__4"]},"data-mapper_cloned_2_cloned":{"name":"customerMapper","type":"DataMapping","datasource":"filter__1","includeExistingColumns":false,"toSchema":[{"fieldName":"account_id","valueExpression":"account_id"},{"fieldName":"created_at","valueExpression":"current_timestamp()"},{"fieldName":"id","valueExpression":"id"},{"fieldName":"status","valueExpression":"bill_status"},{"fieldName":"latest_bill_date","valueExpression":"bill_date"},{"fieldName":"latest_bill_id","valueExpression":"bill_id"},{"fieldName":"last_synced_at","valueExpression":"date_format(current_timestamp(), \"yyyy-MM-dd\\'T\\'HH:mm:ss\")"}],"materialization_strategy":"NONE","fail_on_error":true,"additionalData":{"isGlossaryAssisted":false,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualMappings":[{"id":"mapping-1781114585148","newFieldName":"account_id","mappingType":"sourceColumn","value":"account_id"},{"id":"mapping-1781263224943","newFieldName":"created_at","mappingType":"expression","value":"current_timestamp()"},{"id":"mapping-1783673335306","newFieldName":"id","mappingType":"sourceColumn","value":"id"},{"id":"mapping-1783673699389","newFieldName":"status","mappingType":"sourceColumn","value":"bill_status"},{"id":"mapping-1783673863631","newFieldName":"latest_bill_date","mappingType":"sourceColumn","value":"bill_date"},{"id":"mapping-1783673880596","newFieldName":"latest_bill_id","mappingType":"sourceColumn","value":"bill_id"},{"id":"mapping-1783673951524","newFieldName":"last_synced_at","mappingType":"expression","value":"date_format(current_timestamp(), \"yyyy-MM-dd'T'HH:mm:ss\")"}],"manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[],"after":"","totalSourceTerms":0},"isDefault":false,"connectedComponents":["filter__1"]},"code-transform__0":{"name":"customer_code_transform","type":"CodeTransform","language":"python","datasource":"customerMapper","code":"try:\n # input dataframe will be connected component output as {{datasource}}_df\n # add processing logic here and create output as {{name}}_df\n {{name}}_df = spark.sql(\"\"\"\n SELECT *\n FROM (\n SELECT *,\n ROW_NUMBER() OVER (\n PARTITION BY account_id\n ORDER BY latest_bill_date DESC,\n latest_bill_id DESC\n ) AS rn\n FROM {{datasource}}_df\n ) t\n WHERE rn = 1\n \"\"\") # TODO set output dataframe\n\n # --- Logging additions (safe from template-brace conflicts) ---\n output_df = {{name}}_df\n row_count = output_df.count()\n schema_str = output_df.schema.simpleString()\n print(\"Output row count:\", row_count)\n print(\"Schema:\", schema_str)\n output_df.show(5, truncate=False)\n\n {{name}}_df, {{name}}_observer = observe_metrics(\"{{name}}_df\", {{name}}_df)\n {{name}}_df.createOrReplaceTempView(\"{{name}}_df\")\n\n {{name}}_execute_status=\"SUCCESS\"\nexcept Exception as e:\n print(\"ERROR:\", str(e))\n {{name}}_error = e\n log_error(LOGGER, f\"Component {{name}} Failed\", e)\n {{name}}_execute_status=\"ERROR\"\n raise e","materialization_strategy":"NONE","fail_on_error":true,"isDefault":false,"connectedComponents":["customerMapper"]},"data-mapper_cloned_2_cloned_cloned":{"name":"customerLatestBillMapper","type":"DataMapping","datasource":"customer_code_transform","includeExistingColumns":false,"toSchema":[{"fieldName":"account_id","valueExpression":"account_id"},{"fieldName":"id","valueExpression":"id"},{"fieldName":"last_synced_at","valueExpression":"date_format(current_timestamp(), \"yyyy-MM-dd\\'T\\'HH:mm:ss\")"},{"fieldName":"created_at","valueExpression":"COALESCE(created_at, current_timestamp())"},{"fieldName":"status","valueExpression":"status"},{"fieldName":"latest_bill_date","valueExpression":"latest_bill_date"},{"fieldName":"latest_bill_id","valueExpression":"latest_bill_id"}],"materialization_strategy":"NONE","fail_on_error":true,"additionalData":{"isGlossaryAssisted":false,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualMappings":[{"id":"mapping-1781114585148","newFieldName":"account_id","mappingType":"sourceColumn","value":"account_id"},{"id":"mapping-1783673335306","newFieldName":"id","mappingType":"sourceColumn","value":"id"},{"id":"mapping-1783673951524","newFieldName":"last_synced_at","mappingType":"expression","value":"date_format(current_timestamp(), \"yyyy-MM-dd'T'HH:mm:ss\")"},{"id":"mapping-1783674308922","newFieldName":"created_at","mappingType":"expression","value":"COALESCE(created_at, current_timestamp())"},{"id":"mapping-1783674693520","newFieldName":"status","mappingType":"sourceColumn","value":"status"},{"id":"mapping-1783674703908","newFieldName":"latest_bill_date","mappingType":"sourceColumn","value":"latest_bill_date"},{"id":"mapping-1783674715428","newFieldName":"latest_bill_id","mappingType":"sourceColumn","value":"latest_bill_id"}],"manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[],"after":"","totalSourceTerms":0},"isDefault":false,"connectedComponents":["customer_code_transform"]},"code-transform_cloned":{"name":"CheckpointOutput","type":"CodeTransform","language":"python","datasource":"customerLatestBillMapper","code":"CheckpointOutput_df = {{datasource}}_df.localCheckpoint()\n{{datasource}}_df.persist()\nCheckpointOutput_df.createOrReplaceTempView(\"CheckpointOutput_df\")","materialization_strategy":"NONE","fail_on_error":true,"isDefault":false,"connectedComponents":["customerLatestBillMapper"]},"data-writer_cloned_cloned":{"name":"customer_data_writer","type":"SparkWriter","format":"iceberg","mode":"merge","datasource":"CheckpointOutput","typeLabel":"Spark","credentials":{"accessKey":"S3_ACCESS_KEY","secretKey":"S3_SECRET_KEY"},"iceberg_catalog":"dremio","region":"us-west-1","table_name":"customer","unique_key":["account_id"],"isDefault":false,"connectedComponents":["CheckpointOutput"]}}},"blocks":[{"name":"readCustomers","type":"SparkReader","options":{},"columns":[],"typeLabel":"Spark","datasetName":"","isDefault":false,"connectedComponents":[],"format":"iceberg","credentials":{"accessKey":"S3_ACCESS_KEY","secretKey":"S3_SECRET_KEY"},"iceberg_catalog":"dremio","region":"us-west-1","table_name":"customer"},{"name":"getBills","type":"RESTInvoke","options":{},"datasource":"filterActiveCustomers","url":"https://fw-gateway:8200/fw-notification/outbound-message-config/publish","method":"POST","headers":{"Content-Type":{"value":"application/json","secret":null},"Content-type":{"value":"application/json","secret":null},"api-key":{"value":null,"secret":"OCULAR_API_KEY"},"x-tenantCode":{"value":"UTILITIES","secret":null}},"bodyTemplate":"{\n \"outMsgConfigCode\": \"EXP_ACCOUNT_BILL_HISTORY\",\n \"msgData\": {\n \"accountId\": \"{{account_id}}\",\n \"numberOfMonthPast\": \"24\"\n }\n}","materialization_strategy":"NONE","fail_on_error":true,"isDefault":false,"connectedComponents":["filterActiveCustomers"]},{"name":"filterActiveCustomers","type":"Filter","options":{},"datasource":"readCustomers","condition":"TRIM(UPPER(status)) = \\'ACTIVE\\'","materialization_strategy":"NONE","fail_on_error":true,"isDefault":false,"connectedComponents":["readCustomers"]},{"name":"readLatestBillIds","type":"SparkReader","options":{},"typeLabel":"Spark","columns":[],"isDefault":false,"connectedComponents":[],"format":"iceberg","credentials":{"accessKey":"S3_ACCESS_KEY","secretKey":"S3_SECRET_KEY"},"iceberg_catalog":"dremio","region":"us-west-1","table_name":"bills"},{"name":"MapLatestBill","type":"DataMapping","options":{},"datasource":"getBills","includeExistingColumns":false,"toSchema":[{"fieldName":"accounts","valueExpression":"from_json(\r\n get_json_object(response_body, \\'$.data\\'),\r\n \\'struct<\r\n accountId:string,\r\n numberOfMonthPast:string,\r\n output:struct<\r\n bills:array<struct<\r\n billId:string,\r\n billStatus:string,\r\n billStatusName:string,\r\n billDate:string,\r\n completionDttm:string,\r\n dueDate:string,\r\n amount:string\r\n >>\r\n >\r\n >\\'\r\n)"},{"fieldName":"account_id","valueExpression":"account_id"}],"materialization_strategy":"NONE","fail_on_error":true,"additionalData":{"isGlossaryAssisted":false,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualMappings":[{"id":"mapping-1782329233261","newFieldName":"accounts","mappingType":"expression","value":"from_json(\r\n get_json_object(response_body, '$.data'),\r\n 'struct<\r\n accountId:string,\r\n numberOfMonthPast:string,\r\n output:struct<\r\n bills:array<struct<\r\n billId:string,\r\n billStatus:string,\r\n billStatusName:string,\r\n billDate:string,\r\n completionDttm:string,\r\n dueDate:string,\r\n amount:string\r\n >>\r\n >\r\n >'\r\n)"},{"id":"mapping-1782329846646","newFieldName":"account_id","mappingType":"sourceColumn","value":"account_id"}],"manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[],"after":"","totalSourceTerms":0},"isDefault":false,"connectedComponents":["getBills"]},{"name":"data_join__0","type":"RelationalJoin","options":{},"dropDuplicatedColumns":true,"baseData":"data_mapper__5","joinOrder":[{"with":"readLatestBillIds","joinColumns":[{"account_Id":"account_Id"},{"mapper_bill_id":"bill_id"}],"how":"left outer"}],"isDefault":false,"connectedComponents":["readLatestBillIds","data_mapper__5"]},{"name":"filter__1","type":"Filter","options":{},"datasource":"data_join__0","condition":"bill_id IS NULL OR mapper_bill_id <> bill_id","materialization_strategy":"NONE","fail_on_error":true,"isDefault":false,"connectedComponents":["data_join__0"]},{"name":"BillWriterMapper","type":"DataMapping","options":{},"datasource":"filter__1","includeExistingColumns":false,"toSchema":[{"fieldName":"account_id","valueExpression":"account_id"},{"fieldName":"bill_id","valueExpression":"mapper_bill_id"},{"fieldName":"bill_date","valueExpression":"bill_date"},{"fieldName":"bill_status","valueExpression":"bill_status"},{"fieldName":"due_date","valueExpression":"due_date"},{"fieldName":"created_at","valueExpression":"current_timestamp()"},{"fieldName":"id","valueExpression":"id"},{"fieldName":"amount","valueExpression":"amount"},{"fieldName":"amount_value","valueExpression":"amount_value"},{"fieldName":"bill_status_name","valueExpression":"bill_status_name"},{"fieldName":"completion_dttm","valueExpression":"completion_dttm"}],"materialization_strategy":"NONE","fail_on_error":true,"additionalData":{"isGlossaryAssisted":false,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualMappings":[{"id":"mapping-1781114585148","newFieldName":"account_id","mappingType":"sourceColumn","value":"account_id"},{"id":"mapping-1782331737547","newFieldName":"bill_id","mappingType":"sourceColumn","value":"mapper_bill_id"},{"id":"mapping-1782331770716","newFieldName":"bill_date","mappingType":"sourceColumn","value":"bill_date"},{"id":"mapping-1782331787498","newFieldName":"bill_status","mappingType":"sourceColumn","value":"bill_status"},{"id":"mapping-1782331808058","newFieldName":"due_date","mappingType":"sourceColumn","value":"due_date"},{"id":"mapping-1782331832635","newFieldName":"created_at","mappingType":"expression","value":"current_timestamp()"},{"id":"mapping-1782331863479","newFieldName":"id","mappingType":"sourceColumn","value":"id"},{"id":"mapping-1782331899605","newFieldName":"amount","mappingType":"sourceColumn","value":"amount"},{"id":"mapping-1783665970895","newFieldName":"amount_value","mappingType":"sourceColumn","value":"amount_value"},{"id":"mapping-1783665993133","newFieldName":"bill_status_name","mappingType":"sourceColumn","value":"bill_status_name"},{"id":"mapping-1784539493167","newFieldName":"completion_dttm","mappingType":"sourceColumn","value":"completion_dttm"}],"manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[],"after":"","totalSourceTerms":0},"isDefault":false,"connectedComponents":["filter__1"]},{"name":"data_writer__1","type":"SparkWriter","options":{},"format":"iceberg","mode":"append","datasource":"BillWriterMapper","typeLabel":"Spark","credentials":{"accessKey":"S3_ACCESS_KEY","secretKey":"S3_SECRET_KEY"},"iceberg_catalog":"dremio","region":"us-west-1","table_name":"bills","isDefault":false,"connectedComponents":["BillWriterMapper","forecast_insight_data_mapper"]},{"name":"forecast_insight_code_transform","type":"CodeTransform","options":{},"language":"python","datasource":"data_mapper__5","code":"try:\n\n import builtins\n import json as json_lib\n import traceback\n from datetime import datetime, timedelta\n import numpy as np\n\n # =====================================================\n # CHECK PROPHET\n # =====================================================\n\n try:\n from prophet import Prophet\n PROPHET_AVAILABLE = False\n # print(\"Prophet installed\")\n except Exception as e:\n PROPHET_AVAILABLE = True\n # print(f\"Prophet not available: {e}\")\n\n # =====================================================\n # THRESHOLDS (mirrors BillForecastService class constants)\n # =====================================================\n\n HIGH_USAGE_THRESHOLD = 15.0 # % increase -> high_usage\n DROP_THRESHOLD = -15.0 # % decrease -> drop_detected\n MIN_BILLS_FOR_FILTERING = 3 # minimum bills to apply IQR outlier filtering\n TREND_DAMPEN = 0.5 # apply only 50% of observed MoM change in fallback\n\n # =====================================================\n # READ SOURCE\n # =====================================================\n\n source_df = {{datasource}}_df\n # print(\"Input schema:\")\n source_df.printSchema()\n\n pdf = (\n source_df\n .select(\n \"account_id\",\n \"mapper_bill_id\",\n \"bill_date\",\n \"amount_value\"\n )\n .toPandas()\n )\n\n\n pdf[\"bill_date\"] = pd.to_datetime(pdf[\"bill_date\"])\n\n # print(\"Input rows =\", len(pdf))\n\n # =====================================================\n # MAPPER FILTER — only keep accounts present in mapper_df\n # =====================================================\n \n mapper_df = BillWriterMapper_df\n \n mapper_pdf = (\n mapper_df\n .select(\"account_id\")\n .toPandas()\n )\n \n if mapper_pdf.empty:\n print(\"Mapper has no data - returning empty output successfully\")\n pdf = pdf.iloc[0:0]\n else:\n mapper_account_ids = set(mapper_pdf[\"account_id\"].dropna().unique())\n print(\"Mapper account count =\", len(mapper_account_ids))\n \n before_count = len(pdf)\n pdf = pdf[pdf[\"account_id\"].isin(mapper_account_ids)].reset_index(drop=True)\n print(f\"Filtered source rows by mapper: {before_count} -> {len(pdf)}\")\n \n print(\"Input rows after mapper filter =\", len(pdf))\n\n\n output_rows = []\n\n # =====================================================\n # HELPER FUNCTIONS — outlier filtering & weighting\n # =====================================================\n\n def filter_outliers(amounts):\n \"\"\"Remove outliers using IQR method. Returns filtered list (at least 2 values kept).\"\"\"\n if len(amounts) < 3:\n return amounts\n\n sorted_vals = sorted(amounts)\n n = len(sorted_vals)\n q1 = sorted_vals[n // 4]\n q3 = sorted_vals[(3 * n) // 4]\n iqr = q3 - q1\n\n # Use 1.5x IQR rule; if IQR is 0, fall back to median +/- band\n if iqr > 0:\n lower_bound = q1 - 1.5 * iqr\n upper_bound = q3 + 1.5 * iqr\n else:\n median = sorted_vals[n // 2]\n lower_bound = median * 0.2\n upper_bound = median * 3.0\n\n filtered = [a for a in amounts if lower_bound <= a <= upper_bound]\n\n # Always keep at least the 2 most recent values\n if len(filtered) < 2:\n filtered = amounts[-2:]\n\n return filtered\n\n def exponential_weights(n, decay=0.5):\n \"\"\"Generate exponential decay weights - most recent gets highest weight.\n\n Example with n=3, decay=0.5: [0.25, 0.5, 1.0] -> normalized to [0.143, 0.286, 0.571]\n \"\"\"\n raw = [decay ** (n - 1 - i) for i in range(n)]\n total = builtins.sum(raw)\n return [w / total for w in raw]\n\n def detect_anomalies(amounts, threshold=2.0):\n \"\"\"Detect anomalies using Z-score method.\"\"\"\n if len(amounts) < 3:\n return [False] * len(amounts)\n\n mean_val = np.mean(amounts)\n std_val = np.std(amounts)\n\n if std_val == 0:\n return [False] * len(amounts)\n\n z_scores = [(x - mean_val) / std_val for x in amounts]\n return [bool(abs(z) > threshold) for z in z_scores]\n\n def calculate_trend_slope(amounts):\n \"\"\"Calculate normalized trend slope using linear regression (% change per period).\"\"\"\n if len(amounts) < 2:\n return 0.0\n\n x = np.arange(len(amounts))\n y = np.array(amounts)\n\n n = len(x)\n denom = (n * np.sum(x ** 2) - np.sum(x) ** 2)\n if denom == 0:\n return 0.0\n slope = (n * np.sum(x * y) - np.sum(x) * np.sum(y)) / denom\n\n mean_val = np.mean(amounts)\n if mean_val > 0:\n return (slope / mean_val) * 100\n return 0.0\n\n # =====================================================\n # HELPER FUNCTIONS — forecasting\n # =====================================================\n\n def prophet_forecast(prophet_df, periods=3):\n \"\"\"Run Prophet forecast for the next `periods` months. Raises on failure\n so the caller can fall back to fallback_forecast_weighted.\"\"\"\n model = Prophet(\n yearly_seasonality=True,\n weekly_seasonality=False,\n daily_seasonality=False,\n interval_width=0.80\n )\n\n model.fit(prophet_df)\n\n future = model.make_future_dataframe(\n periods=periods,\n freq=\"M\"\n )\n\n pred = model.predict(future)\n\n forecast_df = pred[\n pred[\"ds\"] > prophet_df[\"ds\"].max()\n ][[\n \"ds\",\n \"yhat\",\n \"yhat_lower\",\n \"yhat_upper\"\n ]].copy()\n\n forecast_df[\"yhat\"] = forecast_df[\"yhat\"].clip(lower=0).round(2)\n forecast_df[\"yhat_lower\"] = forecast_df[\"yhat_lower\"].clip(lower=0).round(2)\n forecast_df[\"yhat_upper\"] = forecast_df[\"yhat_upper\"].round(2)\n\n return forecast_df\n\n def fallback_forecast_weighted(prophet_df, periods=3):\n \"\"\"\n Weighted-average fallback with dampened trend (used when Prophet is\n unavailable, fails, or there are fewer than 4 data points).\n\n Steps:\n 1. Outlier filtering (IQR method) when >= MIN_BILLS_FOR_FILTERING bills.\n 2. Exponential decay weighting (most recent bill weighted highest).\n 3. Dampened month-over-month trend projection (50% of observed rate),\n with a seasonal override when a same-calendar-month average exists.\n\n Confidence interval: +/-15% around the forecast value.\n \"\"\"\n amounts_all = [float(v) for v in prophet_df[\"y\"].values]\n dates_all = list(prophet_df[\"ds\"].values)\n\n amounts = [a for a in amounts_all if a > 0]\n if not amounts:\n return pd.DataFrame(columns=[\"ds\", \"yhat\", \"yhat_lower\", \"yhat_upper\"])\n\n # Step 1: outlier filtering\n if len(amounts) >= MIN_BILLS_FOR_FILTERING:\n clean_amounts = filter_outliers(amounts)\n else:\n clean_amounts = amounts\n\n # Seasonal map: month-of-year -> list of historical amounts in that month\n monthly_map = {}\n for d, a in zip(dates_all, amounts_all):\n if a <= 0:\n continue\n month = pd.Timestamp(d).month\n monthly_map.setdefault(month, []).append(a)\n\n last_date = prophet_df[\"ds\"].max()\n\n # Step 2: exponential decay weighted average on clean data\n weights = exponential_weights(len(clean_amounts))\n weighted_avg = builtins.sum(a * w for a, w in zip(clean_amounts, weights))\n\n # Step 3: dampened month-over-month trend\n if len(clean_amounts) >= 2:\n mom_changes = []\n for j in range(1, len(clean_amounts)):\n if clean_amounts[j - 1] > 0:\n mom_changes.append(\n (clean_amounts[j] - clean_amounts[j - 1]) / clean_amounts[j - 1]\n )\n avg_mom = (builtins.sum(mom_changes) / len(mom_changes)) if mom_changes else 0.0\n dampened_mom = avg_mom * TREND_DAMPEN\n else:\n dampened_mom = 0.0\n\n rows = []\n base_val = weighted_avg\n\n for i in range(1, periods + 1):\n future_dt = last_date + pd.DateOffset(months=i)\n future_month = future_dt.month\n\n if future_month in monthly_map and monthly_map[future_month]:\n seasonal_avg = builtins.sum(monthly_map[future_month]) / len(monthly_map[future_month])\n predicted_value = seasonal_avg\n else:\n predicted_value = builtins.max(0.0, base_val * (1 + dampened_mom) ** i)\n\n lower_bound = builtins.max(0.0, predicted_value * 0.85)\n upper_bound = predicted_value * 1.15\n\n rows.append({\n \"ds\": future_dt,\n \"yhat\": round(predicted_value, 2),\n \"yhat_lower\": round(lower_bound, 2),\n \"yhat_upper\": round(upper_bound, 2)\n })\n\n # print(\n # f\"Fallback forecast: {len(amounts)} bills -> {len(clean_amounts)} clean -> \"\n # f\"base ${weighted_avg:.2f}, dampened MoM {dampened_mom * 100:.1f}%\"\n # )\n\n return pd.DataFrame(rows)\n\n # =====================================================\n # HELPER FUNCTIONS — classification, severity, explanation\n # =====================================================\n\n def classify_type(recent_amounts, forecast_amounts):\n \"\"\"Classify insight type based on % change between recent avg and forecast avg.\"\"\"\n if not recent_amounts or not forecast_amounts:\n return \"stable_usage\"\n\n recent_avg = builtins.sum(recent_amounts) / len(recent_amounts)\n forecast_avg = builtins.sum(forecast_amounts) / len(forecast_amounts)\n\n if recent_avg == 0:\n return \"stable_usage\"\n\n pct_change = ((forecast_avg - recent_avg) / recent_avg) * 100\n\n if pct_change >= HIGH_USAGE_THRESHOLD:\n return \"high_usage\"\n elif pct_change <= DROP_THRESHOLD:\n return \"drop_detected\"\n else:\n return \"stable_usage\"\n\n def compute_severity_score(recent_amounts, forecast_amounts):\n \"\"\"Compute a 1-10 severity score based on magnitude of change.\"\"\"\n if not recent_amounts or not forecast_amounts:\n return 1\n\n recent_avg = builtins.sum(recent_amounts) / len(recent_amounts)\n forecast_avg = builtins.sum(forecast_amounts) / len(forecast_amounts)\n\n if recent_avg == 0:\n return 1\n\n pct_change = abs(((forecast_avg - recent_avg) / recent_avg) * 100)\n return builtins.min(10, builtins.max(1, int(pct_change / 10) + 1))\n\n def generate_explanation_template(recent_amounts, forecast_amounts, insight_type, bill_count):\n \"\"\"Template-based alert/message/explanation for the Rank 1 forecast insight.\"\"\"\n recent_avg = builtins.sum(recent_amounts) / len(recent_amounts) if recent_amounts else 0\n forecast_avg = builtins.sum(forecast_amounts) / len(forecast_amounts) if forecast_amounts else 0\n\n if recent_avg > 0:\n pct_change = ((forecast_avg - recent_avg) / recent_avg) * 100\n else:\n pct_change = 0\n\n direction = \"increase\" if pct_change > 0 else \"decrease\"\n abs_pct = abs(pct_change)\n\n type_labels = {\n \"high_usage\": \"High Usage Expected\",\n \"drop_detected\": \"Bill Drop Detected\",\n \"stable_usage\": \"Stable Billing Pattern\",\n }\n\n alert = type_labels.get(insight_type, \"Bill Forecast\")\n message = f\"Forecasted bills show a {abs_pct:.1f}% {direction} over the next 3 months.\"\n explanation = (\n f\"Based on the last {bill_count} months of billing data, \"\n f\"the average recent bill is ${recent_avg:.2f} and the forecasted average is ${forecast_avg:.2f}. \"\n f\"This represents a {abs_pct:.1f}% {direction} \"\n f\"(${abs(forecast_avg - recent_avg):.2f} difference).\"\n )\n\n return alert, message, explanation\n\n # =====================================================\n # HELPER FUNCTIONS — graphs & considered bills\n # =====================================================\n\n def build_bar_graph(labels, data, label, color=\"rgba(75,192,192,0.6)\"):\n return {\n \"type\": \"bar\",\n \"labels\": labels,\n \"datasets\": [{\n \"label\": label,\n \"data\": data,\n \"backgroundColor\": color\n }]\n }\n\n def build_forecast_graph(forecast_labels, forecast_values, forecast_lower, forecast_upper):\n return {\n \"type\": \"bar\",\n \"labels\": forecast_labels,\n \"datasets\": [\n {\n \"label\": \"Forecasted Bills\",\n \"data\": forecast_values,\n \"borderColor\": \"rgb(75,102,192)\",\n \"backgroundColor\": \"rgba(75,192,192,0.2)\",\n \"fill\": True\n },\n {\n \"label\": \"Confidence Lower\",\n \"data\": forecast_lower,\n \"borderColor\": \"rgba(75,192,192,0.3)\",\n \"backgroundColor\": \"transparent\",\n \"borderDash\": [5, 5],\n \"fill\": False\n },\n {\n \"label\": \"Confidence Upper\",\n \"data\": forecast_upper,\n \"borderColor\": \"rgba(75,192,192,0.3)\",\n \"backgroundColor\": \"transparent\",\n \"borderDash\": [5, 5],\n \"fill\": False\n }\n ]\n }\n\n def build_considered_bills(rows_df, anomaly_flags=None):\n \"\"\"Build the consideredBills list from a pandas slice of bill rows.\"\"\"\n considered = []\n for i, (_, r) in enumerate(rows_df.iterrows()):\n is_anomaly = bool(anomaly_flags[i]) if anomaly_flags and i < len(anomaly_flags) else False\n considered.append({\n \"billID\": str(r[\"mapper_bill_id\"]),\n \"billDate\": r[\"bill_date\"].strftime(\"%Y-%m-%d\"),\n \"billAmount\": round(float(r[\"amount_value\"]), 2),\n \"consumptionValue\": round(float(r[\"amount_value\"]), 2),\n \"consumptionUnit\": \"USD\",\n \"isAnomaly\": is_anomaly\n })\n return considered\n\n # =====================================================\n # INSIGHT BUILDER — Rank 1: Bill Forecast\n # =====================================================\n\n def build_forecast_insight(recent, forecast_df, anomaly_flags):\n \"\"\"\n Rank 1 insight: forecast classification (high_usage / drop_detected /\n stable_usage), with historical graph + forecast graph + consideredBills.\n \"\"\"\n actual_labels = [d.strftime(\"%Y-%m\") for d in recent[\"bill_date\"]]\n actual_amounts = [round(float(x), 2) for x in recent[\"amount_value\"]]\n\n forecast_labels = [d.strftime(\"%Y-%m\") for d in forecast_df[\"ds\"]]\n forecast_values = [round(float(x), 2) for x in forecast_df[\"yhat\"]]\n forecast_lower = [round(float(x), 2) for x in forecast_df[\"yhat_lower\"]]\n forecast_upper = [round(float(x), 2) for x in forecast_df[\"yhat_upper\"]]\n\n insight_type = classify_type(actual_amounts, forecast_values)\n severity = compute_severity_score(actual_amounts, forecast_values)\n alert, message, explanation = generate_explanation_template(\n actual_amounts, forecast_values, insight_type, len(recent)\n )\n\n considered_bills = build_considered_bills(recent, anomaly_flags)\n\n actual_graph = build_bar_graph(actual_labels, actual_amounts, \"Bills USD\")\n forecast_graph = build_forecast_graph(forecast_labels, forecast_values, forecast_lower, forecast_upper)\n\n insight = {\n \"rank\": 1,\n \"alert\": alert,\n \"message\": message,\n \"explanation\": explanation,\n \"severityScore\": severity,\n \"consideredBills\": considered_bills,\n \"graph\": actual_graph,\n \"forecastGraph\": forecast_graph,\n \"type\": insight_type\n }\n\n return insight, forecast_labels, forecast_values, forecast_lower, forecast_upper\n\n # =====================================================\n # INSIGHT BUILDER — Rank 2: Trend Summary (last 3 months)\n # =====================================================\n\n def build_trend_insight(recent, anomaly_flags, forecast_labels, forecast_values, forecast_lower, forecast_upper):\n \"\"\"\n Rank 2 insight: month-over-month trend pattern across the last 3 bills.\n\n Patterns:\n declining_trend - both MoM changes < -20%\n increasing_trend - both MoM changes > +20%\n spike_resolved - oldest month 30%+ higher, bills dropped since\n mid_spike - middle month 30%+ higher than neighbors\n recent_spike - most recent month jumped 30%+\n stable_trend - all within 20% of 3-month average\n\n Returns None if fewer than 3 bills or no clear pattern.\n \"\"\"\n if len(recent) < 3:\n return None\n\n amounts = [round(float(x), 2) for x in recent[\"amount_value\"]]\n dates = [d.strftime(\"%Y-%m\") for d in recent[\"bill_date\"]]\n month_names = [d.strftime(\"%B %Y\") for d in recent[\"bill_date\"]]\n\n a0, a1, a2 = amounts # oldest -> newest\n\n def pct(old, new):\n return ((new - old) / old * 100) if old != 0 else 0\n\n chg_1 = pct(a0, a1)\n chg_2 = pct(a1, a2)\n total_chg = pct(a0, a2)\n\n peak_idx = amounts.index(builtins.max(amounts))\n\n alert = \"\"\n message = \"\"\n explanation = \"\"\n trend_type = \"stable_trend\"\n\n if chg_1 < -20 and chg_2 < -20:\n trend_type = \"declining_trend\"\n alert = \"Bills Declining Steadily\"\n message = (\n f\"Your bill has dropped {abs(total_chg):.0f}% over the last 3 months \"\n f\"- from ${a0:,.2f} in {month_names[0]} to ${a2:,.2f} in {month_names[2]}.\"\n )\n explanation = (\n f\"{month_names[0]}: ${a0:,.2f} -> {month_names[1]}: ${a1:,.2f} ({chg_1:+.1f}%) -> \"\n f\"{month_names[2]}: ${a2:,.2f} ({chg_2:+.1f}%). \"\n f\"This is a consistent downward trend that may continue.\"\n )\n\n elif chg_1 > 20 and chg_2 > 20:\n trend_type = \"increasing_trend\"\n alert = \"Bills Increasing Steadily\"\n message = (\n f\"Your bill has risen {abs(total_chg):.0f}% over the last 3 months \"\n f\"- from ${a0:,.2f} in {month_names[0]} to ${a2:,.2f} in {month_names[2]}.\"\n )\n explanation = (\n f\"{month_names[0]}: ${a0:,.2f} -> {month_names[1]}: ${a1:,.2f} ({chg_1:+.1f}%) -> \"\n f\"{month_names[2]}: ${a2:,.2f} ({chg_2:+.1f}%). \"\n f\"Your usage has been climbing; consider reviewing recent activity.\"\n )\n\n elif peak_idx == 0 and abs(total_chg) > 30:\n trend_type = \"spike_resolved\"\n alert = \"Recent Bill Spike Has Resolved\"\n message = (\n f\"Your bill was ${a0:,.2f} in {month_names[0]} but has since dropped \"\n f\"to ${a2:,.2f} in {month_names[2]}.\"\n )\n explanation = (\n f\"The {month_names[0]} bill (${a0:,.2f}) was significantly higher than the recent \"\n f\"{month_names[1]} (${a1:,.2f}) and {month_names[2]} (${a2:,.2f}). \"\n f\"This suggests the spike was a one-time event and bills are normalizing.\"\n )\n\n elif peak_idx == 1 and pct(a1, a0) < -30 and pct(a1, a2) < -30:\n trend_type = \"mid_spike\"\n alert = f\"Bill Spike in {month_names[1]}\"\n message = (\n f\"Your {month_names[1]} bill spiked to ${a1:,.2f} but has returned \"\n f\"to ${a2:,.2f} in {month_names[2]}.\"\n )\n explanation = (\n f\"{month_names[0]}: ${a0:,.2f} -> {month_names[1]}: ${a1:,.2f} \"\n f\"(spike of {pct(a0, a1):+.1f}%) -> \"\n f\"{month_names[2]}: ${a2:,.2f} (back to {pct(a1, a2):+.1f}%). \"\n f\"The {month_names[1]} spike appears to be an anomaly.\"\n )\n\n elif peak_idx == 2 and pct(a1, a2) > 30:\n trend_type = \"recent_spike\"\n alert = \"Recent Bill Spike\"\n message = (\n f\"Your latest bill in {month_names[2]} jumped to ${a2:,.2f} \"\n f\"- up {pct(a1, a2):.0f}% from {month_names[1]} (${a1:,.2f}).\"\n )\n explanation = (\n f\"{month_names[0]}: ${a0:,.2f} -> {month_names[1]}: ${a1:,.2f} -> \"\n f\"{month_names[2]}: ${a2:,.2f} ({pct(a1, a2):+.1f}%). \"\n f\"This recent increase is worth monitoring.\"\n )\n\n else:\n avg_3 = builtins.sum(amounts) / 3\n max_dev = builtins.max(abs(a - avg_3) / avg_3 * 100 for a in amounts) if avg_3 > 0 else 0\n if max_dev < 20:\n trend_type = \"stable_trend\"\n alert = \"Bills Are Stable\"\n message = f\"Your bills have been consistent over the last 3 months, averaging ${avg_3:,.2f}.\"\n explanation = (\n f\"{month_names[0]}: ${a0:,.2f}, {month_names[1]}: ${a1:,.2f}, {month_names[2]}: ${a2:,.2f}. \"\n f\"Variation is within normal range.\"\n )\n else:\n return None\n\n severity = builtins.min(10, builtins.max(1, int(abs(total_chg) / 15) + 1))\n\n considered_bills = build_considered_bills(recent, anomaly_flags)\n\n trend_graph = {\n \"type\": \"bar\",\n \"labels\": dates,\n \"datasets\": [{\n \"label\": \"Monthly Bills USD\",\n \"data\": amounts,\n \"backgroundColor\": [\n \"rgba(255,99,132,0.6)\" if i == peak_idx else \"rgba(75,192,192,0.6)\"\n for i in range(3)\n ]\n }]\n }\n\n forecast_graph = build_forecast_graph(forecast_labels, forecast_values, forecast_lower, forecast_upper)\n\n return {\n \"rank\": 2,\n \"alert\": alert,\n \"message\": message,\n \"explanation\": explanation,\n \"severityScore\": severity,\n \"consideredBills\": considered_bills,\n \"graph\": trend_graph,\n \"forecastGraph\": forecast_graph,\n \"type\": trend_type\n }\n\n # =====================================================\n # INSIGHT BUILDER — Rank 3: Year-over-Year Comparison\n # =====================================================\n\n def build_yoy_insight(acct_df, anomaly_flags_full):\n \"\"\"\n Rank 3 insight: compares the most recent 3 months against the same 3\n calendar months a year ago. Requires >= 6 months of history overall,\n and requires that the same-month-last-year data actually exists.\n Returns None if insufficient data.\n \"\"\"\n if len(acct_df) < 6:\n return None\n\n sorted_df = acct_df.sort_values(\"bill_date\").reset_index(drop=True)\n recent_3 = sorted_df.tail(3)\n recent_dates = [d for d in recent_3[\"bill_date\"]]\n\n yoy_targets = [(d.year - 1, d.month) for d in recent_dates]\n\n yoy_rows = sorted_df[\n sorted_df[\"bill_date\"].apply(lambda d: (d.year, d.month) in yoy_targets)\n ]\n\n if len(yoy_rows) < len(yoy_targets):\n return None\n\n yoy_rows = yoy_rows.sort_values(\"bill_date\").tail(len(yoy_targets))\n\n recent_amounts = [round(float(x), 2) for x in recent_3[\"amount_value\"]]\n yoy_amounts = [round(float(x), 2) for x in yoy_rows[\"amount_value\"]]\n\n recent_avg = builtins.sum(recent_amounts) / len(recent_amounts)\n yoy_avg = builtins.sum(yoy_amounts) / len(yoy_amounts)\n pct_change = ((recent_avg - yoy_avg) / yoy_avg * 100) if yoy_avg != 0 else 0\n\n direction = \"increased\" if pct_change > 0 else \"decreased\"\n if pct_change > HIGH_USAGE_THRESHOLD:\n insight_type = \"high_usage\"\n elif pct_change < DROP_THRESHOLD:\n insight_type = \"drop_detected\"\n else:\n insight_type = \"stable_usage\"\n\n severity = builtins.min(10, builtins.max(1, int(abs(pct_change) / 10) + 1))\n\n # anomaly flags computed over the full account history align by position\n recent_anomalies = anomaly_flags_full[-3:] if len(anomaly_flags_full) >= 3 else [False] * 3\n considered_bills = build_considered_bills(recent_3, recent_anomalies) + build_considered_bills(yoy_rows, None)\n\n yoy_graph = build_bar_graph(\n [d.strftime(\"%Y-%m\") for d in yoy_rows[\"bill_date\"]],\n yoy_amounts,\n \"Last Year Bills USD\",\n color=\"rgba(153,102,255,0.6)\"\n )\n current_graph = build_bar_graph(\n [d.strftime(\"%Y-%m\") for d in recent_3[\"bill_date\"]],\n recent_amounts,\n \"Current Year Bills USD\",\n color=\"rgba(75,192,192,0.6)\"\n )\n\n return {\n \"rank\": 3,\n \"alert\": f\"Year-over-Year Bill {direction.capitalize()}\",\n \"message\": f\"Your bills have {direction} by {abs(pct_change):.1f}% compared to the same period last year.\",\n \"explanation\": (\n f\"Average bill for the recent 3 months: ${recent_avg:.2f}. \"\n f\"Average bill for the same 3 months last year: ${yoy_avg:.2f}. \"\n f\"That is a {abs(pct_change):.1f}% {direction}.\"\n ),\n \"severityScore\": severity,\n \"consideredBills\": considered_bills,\n \"graph\": yoy_graph,\n \"forecastGraph\": current_graph,\n \"type\": insight_type\n }\n\n # =====================================================\n # FORECAST ACCURACY — best-effort in-sample backtest\n # =====================================================\n\n def compute_forecast_accuracy_backtest(acct_df):\n \"\"\"\n Best-effort forecast accuracy, computed entirely from the bills already\n present in this dataframe (no external previous-forecast input available).\n\n Approach: hold out the most recent actual bill, forecast 1 month ahead\n using only the months before it (same Prophet/fallback logic as the\n live forecast), then compare that 1-month-ahead prediction against the\n real bill that came in. This approximates \"how accurate was last\n month's forecast\" without needing a stored previous forecast.\n\n Requires >= 4 bills (3 to forecast from + 1 actual to validate against).\n Returns None if not enough data.\n \"\"\"\n sorted_df = acct_df.sort_values(\"bill_date\").reset_index(drop=True)\n if len(sorted_df) < 4:\n return None\n\n train_df = sorted_df.iloc[:-1]\n actual_row = sorted_df.iloc[-1]\n\n train_prophet_df = train_df.rename(columns={\"bill_date\": \"ds\", \"amount_value\": \"y\"})[[\"ds\", \"y\"]]\n\n try:\n if PROPHET_AVAILABLE and len(train_df) >= 4:\n bt_forecast_df = prophet_forecast(train_prophet_df, periods=1)\n else:\n raise Exception(\"Prophet unavailable or insufficient data for backtest\")\n except Exception:\n bt_forecast_df = fallback_forecast_weighted(train_prophet_df, periods=1)\n\n if bt_forecast_df.empty:\n return None\n\n predicted = float(bt_forecast_df.iloc[0][\"yhat\"])\n actual = float(actual_row[\"amount_value\"])\n\n error_pct = abs(predicted - actual) / actual * 100 if actual != 0 else 0\n accuracy = builtins.max(0, 100 - error_pct)\n\n return {\n \"method\": \"in_sample_backtest\",\n \"validatedMonth\": actual_row[\"bill_date\"].strftime(\"%Y-%m\"),\n \"predicted\": round(predicted, 2),\n \"actual\": round(actual, 2),\n \"accuracyPct\": round(accuracy, 1)\n }\n\n # =====================================================\n # PROCESS EACH ACCOUNT\n # =====================================================\n\n for account_id, acct_df in pdf.groupby(\"account_id\"):\n\n try:\n\n # print(f\"\\nProcessing account {account_id}\")\n\n acct_df = acct_df.sort_values(\"bill_date\").reset_index(drop=True)\n\n bill_count = len(acct_df)\n\n # print(\"Bill count =\", bill_count)\n\n if bill_count < 3:\n # print(\"Skipping account - less than 3 bills\")\n continue\n\n # =============================================\n # PROPHET INPUT\n # =============================================\n\n prophet_df = acct_df.rename(\n columns={\n \"bill_date\": \"ds\",\n \"amount_value\": \"y\"\n }\n )[[\"ds\", \"y\"]]\n\n # =============================================\n # FORECAST (Prophet >= 4 points, else weighted fallback)\n # =============================================\n\n try:\n\n if PROPHET_AVAILABLE and bill_count >= 4:\n\n # print(\"Running Prophet\")\n forecast_df = prophet_forecast(prophet_df, periods=3)\n\n else:\n raise Exception(\"Prophet unavailable or insufficient data (<4 points)\")\n\n except Exception as prophet_error:\n\n # print(f\"Prophet failed for {account_id}: {prophet_error}\")\n # print(\"Using weighted-average fallback (outlier filtering + exponential decay + dampened trend)\")\n\n forecast_df = fallback_forecast_weighted(prophet_df, periods=3)\n\n # print(\"Forecast rows =\", len(forecast_df))\n # print(\"forecast_df\", forecast_df)\n\n # =============================================\n # ANOMALY DETECTION (full history, Z-score)\n # =============================================\n\n all_amounts = [round(float(x), 2) for x in acct_df[\"amount_value\"]]\n anomaly_flags_full = detect_anomalies(all_amounts)\n\n # =============================================\n # ACTUAL DATA — recent 3 months\n # =============================================\n\n recent = acct_df.tail(3).reset_index(drop=True)\n recent_anomaly_flags = anomaly_flags_full[-3:] if len(anomaly_flags_full) >= 3 else [False] * len(recent)\n\n # print(\"recent\", recent)\n\n # =============================================\n # TREND SLOPE (informational, kept in output)\n # =============================================\n\n recent_amounts_for_slope = [round(float(x), 2) for x in recent[\"amount_value\"]]\n trend_slope = calculate_trend_slope(recent_amounts_for_slope)\n\n # =============================================\n # RANK 1 — FORECAST INSIGHT\n # =============================================\n\n forecast_insight, forecast_labels, forecast_values, forecast_lower, forecast_upper = (\n build_forecast_insight(recent, forecast_df, recent_anomaly_flags)\n )\n\n insights = [forecast_insight]\n\n # =============================================\n # RANK 2 — TREND SUMMARY\n # =============================================\n\n trend_insight = build_trend_insight(\n recent, recent_anomaly_flags,\n forecast_labels, forecast_values, forecast_lower, forecast_upper\n )\n if trend_insight:\n insights.append(trend_insight)\n\n # =============================================\n # RANK 3 — YEAR-OVER-YEAR COMPARISON\n # =============================================\n\n yoy_insight = build_yoy_insight(acct_df, anomaly_flags_full)\n if yoy_insight:\n insights.append(yoy_insight)\n\n # =============================================\n # FORECAST ACCURACY — best-effort backtest\n # =============================================\n\n accuracy_data = compute_forecast_accuracy_backtest(acct_df)\n if accuracy_data:\n for ins in insights:\n if ins.get(\"rank\") == 1:\n ins[\"previousForecastAccuracy\"] = accuracy_data\n\n # =============================================\n # FINAL JSON\n # =============================================\n\n response = {\n \"accountId\": account_id,\n \"generatedAt\": datetime.utcnow().isoformat(),\n \"cacheHit\": False,\n \"dataPointsUsed\": len(recent),\n \"nextRefreshDate\": (\n datetime.utcnow() + pd.DateOffset(months=1)\n ).strftime(\"%Y-%m-%d\"),\n \"forecastMonth\": datetime.utcnow().strftime(\"%Y-%m\"),\n \"stale\": False,\n \"forecastMethod\": \"prophet\" if (PROPHET_AVAILABLE and bill_count >= 4) else \"weighted_average_fallback\",\n \"trendSlope\": round(trend_slope, 2),\n \"insights\": insights\n }\n\n output_rows.append(\n Row(\n account_id=str(account_id),\n insight_json=json_lib.dumps(response)\n )\n )\n\n except Exception as account_error:\n\n # print(f\"Account failed {account_id}: {account_error}\")\n\n traceback.print_exc()\n continue\n\n # =====================================================\n # OUTPUT\n # =====================================================\n\n if len(output_rows) > 0:\n\n forecast_insight_code_transform_df = (\n spark.createDataFrame(output_rows)\n )\n\n else:\n\n empty_schema = (\n \"account_id string,\"\n \" insight_json string\"\n )\n\n forecast_insight_code_transform_df = (\n spark.createDataFrame(\n [],\n empty_schema\n )\n )\n\n # print(f\"Generated insights for {len(output_rows)} accounts\")\n\n # forecast_insight_code_transform_df.show(truncate=False)\n\n forecast_insight_code_transform_df.createOrReplaceTempView(\"{{name}}_df\")\n\n forecast_insight_code_transform_execute_status = \"SUCCESS\"\n\nexcept Exception as e:\n\n print(\"Pipeline failed\")\n print(str(e))\n\n forecast_insight_code_transform_execute_status = \"ERROR\"\n\n raise","materialization_strategy":"NONE","fail_on_error":true,"isDefault":false,"connectedComponents":["BillWriterMapper","data_mapper__5"]},{"name":"forecast_insight_data_mapper","type":"DataMapping","options":{},"datasource":"forecast_insight_code_transform","includeExistingColumns":false,"toSchema":[{"fieldName":"account_id","valueExpression":"account_id"},{"fieldName":"id","valueExpression":"uuid()"},{"fieldName":"insight_data","valueExpression":"get_json_object(insight_json, \\'$.insights\\')"},{"fieldName":"forecast_month","valueExpression":"get_json_object(insight_json, \\'$.forecastMonth\\')"},{"fieldName":"type","valueExpression":"get_json_object(insight_json, \\'$.insights[0].type\\')"},{"fieldName":"created_at","valueExpression":"date_format(current_timestamp(), \"yyyy-MM-dd\\'T\\'HH:mm:ss\")"},{"fieldName":"updated_at","valueExpression":"date_format(current_timestamp(), \"yyyy-MM-dd\\'T\\'HH:mm:ss\")"}],"materialization_strategy":"NONE","fail_on_error":true,"additionalData":{"isGlossaryAssisted":false,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualMappings":[{"id":"mapping-1782062869403","newFieldName":"account_id","mappingType":"sourceColumn","value":"account_id"},{"id":"mapping-1782155700948","newFieldName":"id","mappingType":"expression","value":"uuid()"},{"id":"mapping-1782156149735","newFieldName":"insight_data","mappingType":"expression","value":"get_json_object(insight_json, '$.insights')"},{"id":"mapping-1782156157631","newFieldName":"forecast_month","mappingType":"expression","value":"get_json_object(insight_json, '$.forecastMonth')"},{"id":"mapping-1782156764648","newFieldName":"type","mappingType":"expression","value":"get_json_object(insight_json, '$.insights[0].type')"},{"id":"mapping-1783671861263","newFieldName":"created_at","mappingType":"expression","value":"date_format(current_timestamp(), \"yyyy-MM-dd'T'HH:mm:ss\")"},{"id":"mapping-1783671876510","newFieldName":"updated_at","mappingType":"expression","value":"date_format(current_timestamp(), \"yyyy-MM-dd'T'HH:mm:ss\")"}],"manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[],"after":"","totalSourceTerms":0},"isDefault":false,"connectedComponents":["forecast_insight_code_transform"]},{"name":"accountinsights_data_writer","type":"SparkWriter","options":{},"format":"iceberg","mode":"append","datasource":"forecast_insight_data_mapper","typeLabel":"Spark","credentials":{"accessKey":"S3_ACCESS_KEY","secretKey":"S3_SECRET_KEY"},"iceberg_catalog":"dremio","region":"us-west-1","table_name":"insights","isDefault":false,"connectedComponents":["forecast_insight_data_mapper"]},{"name":"data_mapper__3","type":"DataMapping","options":{},"datasource":"MapLatestBill","includeExistingColumns":false,"toSchema":[{"fieldName":"account_id","valueExpression":"account_id"},{"fieldName":"bills","valueExpression":"accounts.output.bills"}],"materialization_strategy":"NONE","fail_on_error":true,"additionalData":{"isGlossaryAssisted":false,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualMappings":[{"id":"mapping-1782329935585","newFieldName":"account_id","mappingType":"sourceColumn","value":"account_id"},{"id":"mapping-1782330182362","newFieldName":"bills","mappingType":"expression","value":"accounts.output.bills"}],"manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[],"after":"","totalSourceTerms":0},"isDefault":false,"connectedComponents":["MapLatestBill"]},{"name":"data_mapper__4","type":"DataMapping","options":{},"datasource":"data_mapper__3","includeExistingColumns":false,"toSchema":[{"fieldName":"account_id","valueExpression":"account_id"},{"fieldName":"final_bills","valueExpression":"explode(bills)"}],"materialization_strategy":"NONE","fail_on_error":true,"additionalData":{"isGlossaryAssisted":false,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualMappings":[{"id":"mapping-1782330221700","newFieldName":"account_id","mappingType":"sourceColumn","value":"account_id"},{"id":"mapping-1782330259011","newFieldName":"final_bills","mappingType":"expression","value":"explode(bills)"}],"manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[],"after":"","totalSourceTerms":0},"isDefault":false,"connectedComponents":["data_mapper__3"]},{"name":"data_mapper__5","type":"DataMapping","options":{},"datasource":"data_mapper__4","includeExistingColumns":false,"toSchema":[{"fieldName":"account_id","valueExpression":"account_id"},{"fieldName":"created_at","valueExpression":"current_timestamp()"},{"fieldName":"id","valueExpression":"uuid()"},{"fieldName":"bill_date","valueExpression":"to_date(final_bills.billDate)"},{"fieldName":"due_date","valueExpression":"to_date(final_bills.dueDate)"},{"fieldName":"bill_status","valueExpression":"final_bills.billStatus"},{"fieldName":"mapper_bill_id","valueExpression":"final_bills.billId"},{"fieldName":"amount_value","valueExpression":"cast(replace(replace(final_bills.amount, \\'$\\', \\'\\'), \\',\\', \\'\\')AS decimal(10, 2))"},{"fieldName":"amount","valueExpression":"final_bills.amount"},{"fieldName":"bill_status_name","valueExpression":"final_bills.billStatusName"},{"fieldName":"completion_dttm","valueExpression":"final_bills.completionDttm"}],"materialization_strategy":"NONE","fail_on_error":true,"additionalData":{"isGlossaryAssisted":false,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualMappings":[{"id":"mapping-1781114585148","newFieldName":"account_id","mappingType":"sourceColumn","value":"account_id"},{"id":"mapping-1781263224943","newFieldName":"created_at","mappingType":"expression","value":"current_timestamp()"},{"id":"mapping-1782190167525","newFieldName":"id","mappingType":"expression","value":"uuid()"},{"id":"mapping-1782331048880","newFieldName":"bill_date","mappingType":"expression","value":"to_date(final_bills.billDate)"},{"id":"mapping-1782331060005","newFieldName":"due_date","mappingType":"expression","value":"to_date(final_bills.dueDate)"},{"id":"mapping-1782331069120","newFieldName":"bill_status","mappingType":"expression","value":"final_bills.billStatus"},{"id":"mapping-1782331518221","newFieldName":"mapper_bill_id","mappingType":"expression","value":"final_bills.billId"},{"id":"mapping-1783664487911","newFieldName":"amount_value","mappingType":"expression","value":"cast(replace(replace(final_bills.amount, '$', ''), ',', '')AS decimal(10, 2))"},{"id":"mapping-1783664496058","newFieldName":"amount","mappingType":"expression","value":"final_bills.amount"},{"id":"mapping-1783664524438","newFieldName":"bill_status_name","mappingType":"expression","value":"final_bills.billStatusName"},{"id":"mapping-1783664548606","newFieldName":"completion_dttm","mappingType":"expression","value":"final_bills.completionDttm"}],"manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[],"after":"","totalSourceTerms":0},"isDefault":false,"connectedComponents":["data_mapper__4"]},{"name":"customerMapper","type":"DataMapping","options":{},"datasource":"filter__1","includeExistingColumns":false,"toSchema":[{"fieldName":"account_id","valueExpression":"account_id"},{"fieldName":"created_at","valueExpression":"current_timestamp()"},{"fieldName":"id","valueExpression":"id"},{"fieldName":"status","valueExpression":"bill_status"},{"fieldName":"latest_bill_date","valueExpression":"bill_date"},{"fieldName":"latest_bill_id","valueExpression":"bill_id"},{"fieldName":"last_synced_at","valueExpression":"date_format(current_timestamp(), \"yyyy-MM-dd\\'T\\'HH:mm:ss\")"}],"materialization_strategy":"NONE","fail_on_error":true,"additionalData":{"isGlossaryAssisted":false,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualMappings":[{"id":"mapping-1781114585148","newFieldName":"account_id","mappingType":"sourceColumn","value":"account_id"},{"id":"mapping-1781263224943","newFieldName":"created_at","mappingType":"expression","value":"current_timestamp()"},{"id":"mapping-1783673335306","newFieldName":"id","mappingType":"sourceColumn","value":"id"},{"id":"mapping-1783673699389","newFieldName":"status","mappingType":"sourceColumn","value":"bill_status"},{"id":"mapping-1783673863631","newFieldName":"latest_bill_date","mappingType":"sourceColumn","value":"bill_date"},{"id":"mapping-1783673880596","newFieldName":"latest_bill_id","mappingType":"sourceColumn","value":"bill_id"},{"id":"mapping-1783673951524","newFieldName":"last_synced_at","mappingType":"expression","value":"date_format(current_timestamp(), \"yyyy-MM-dd'T'HH:mm:ss\")"}],"manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[],"after":"","totalSourceTerms":0},"isDefault":false,"connectedComponents":["filter__1"]},{"name":"customer_code_transform","type":"CodeTransform","options":{},"language":"python","datasource":"customerMapper","code":"try:\n # input dataframe will be connected component output as {{datasource}}_df\n # add processing logic here and create output as {{name}}_df\n {{name}}_df = spark.sql(\"\"\"\n SELECT *\n FROM (\n SELECT *,\n ROW_NUMBER() OVER (\n PARTITION BY account_id\n ORDER BY latest_bill_date DESC,\n latest_bill_id DESC\n ) AS rn\n FROM {{datasource}}_df\n ) t\n WHERE rn = 1\n \"\"\") # TODO set output dataframe\n\n # --- Logging additions (safe from template-brace conflicts) ---\n output_df = {{name}}_df\n row_count = output_df.count()\n schema_str = output_df.schema.simpleString()\n print(\"Output row count:\", row_count)\n print(\"Schema:\", schema_str)\n output_df.show(5, truncate=False)\n\n {{name}}_df, {{name}}_observer = observe_metrics(\"{{name}}_df\", {{name}}_df)\n {{name}}_df.createOrReplaceTempView(\"{{name}}_df\")\n\n {{name}}_execute_status=\"SUCCESS\"\nexcept Exception as e:\n print(\"ERROR:\", str(e))\n {{name}}_error = e\n log_error(LOGGER, f\"Component {{name}} Failed\", e)\n {{name}}_execute_status=\"ERROR\"\n raise e","materialization_strategy":"NONE","fail_on_error":true,"isDefault":false,"connectedComponents":["customerMapper"]},{"name":"customerLatestBillMapper","type":"DataMapping","options":{},"datasource":"customer_code_transform","includeExistingColumns":false,"toSchema":[{"fieldName":"account_id","valueExpression":"account_id"},{"fieldName":"id","valueExpression":"id"},{"fieldName":"last_synced_at","valueExpression":"date_format(current_timestamp(), \"yyyy-MM-dd\\'T\\'HH:mm:ss\")"},{"fieldName":"created_at","valueExpression":"COALESCE(created_at, current_timestamp())"},{"fieldName":"status","valueExpression":"status"},{"fieldName":"latest_bill_date","valueExpression":"latest_bill_date"},{"fieldName":"latest_bill_id","valueExpression":"latest_bill_id"}],"materialization_strategy":"NONE","fail_on_error":true,"additionalData":{"isGlossaryAssisted":false,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualMappings":[{"id":"mapping-1781114585148","newFieldName":"account_id","mappingType":"sourceColumn","value":"account_id"},{"id":"mapping-1783673335306","newFieldName":"id","mappingType":"sourceColumn","value":"id"},{"id":"mapping-1783673951524","newFieldName":"last_synced_at","mappingType":"expression","value":"date_format(current_timestamp(), \"yyyy-MM-dd'T'HH:mm:ss\")"},{"id":"mapping-1783674308922","newFieldName":"created_at","mappingType":"expression","value":"COALESCE(created_at, current_timestamp())"},{"id":"mapping-1783674693520","newFieldName":"status","mappingType":"sourceColumn","value":"status"},{"id":"mapping-1783674703908","newFieldName":"latest_bill_date","mappingType":"sourceColumn","value":"latest_bill_date"},{"id":"mapping-1783674715428","newFieldName":"latest_bill_id","mappingType":"sourceColumn","value":"latest_bill_id"}],"manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[],"after":"","totalSourceTerms":0},"isDefault":false,"connectedComponents":["customer_code_transform"]},{"name":"CheckpointOutput","type":"CodeTransform","options":{},"language":"python","datasource":"customerLatestBillMapper","code":"CheckpointOutput_df = {{datasource}}_df.localCheckpoint()\n{{datasource}}_df.persist()\nCheckpointOutput_df.createOrReplaceTempView(\"CheckpointOutput_df\")","materialization_strategy":"NONE","fail_on_error":true,"isDefault":false,"connectedComponents":["customerLatestBillMapper"]},{"name":"customer_data_writer","type":"SparkWriter","options":{},"format":"iceberg","mode":"merge","datasource":"CheckpointOutput","typeLabel":"Spark","credentials":{"accessKey":"S3_ACCESS_KEY","secretKey":"S3_SECRET_KEY"},"iceberg_catalog":"dremio","region":"us-west-1","table_name":"customer","unique_key":["account_id"],"isDefault":false,"connectedComponents":["CheckpointOutput"]}]}} |