From 505b9cc65426aec4eaef36408ae620d0de7b9b8f Mon Sep 17 00:00:00 2001 From: Max Alexander Date: Thu, 5 Feb 2026 18:12:18 -0500 Subject: [PATCH 1/3] Add logging for files detected and files written --- .../ex_cubic_ingestion/process_incoming.ex | 37 +++++++++++++------ .../lib/py_cubic_ingestion/ingest_incoming.py | 10 +++++ 2 files changed, 36 insertions(+), 11 deletions(-) diff --git a/ex_cubic_ingestion/lib/ex_cubic_ingestion/process_incoming.ex b/ex_cubic_ingestion/lib/ex_cubic_ingestion/process_incoming.ex index 5ee2a73..af0ba94 100644 --- a/ex_cubic_ingestion/lib/ex_cubic_ingestion/process_incoming.ex +++ b/ex_cubic_ingestion/lib/ex_cubic_ingestion/process_incoming.ex @@ -14,6 +14,7 @@ defmodule ExCubicIngestion.ProcessIncoming do alias ExCubicIngestion.Schema.CubicLoad alias ExCubicIngestion.Schema.CubicTable alias ExCubicIngestion.Validators + require Logger @wait_interval_ms 5_000 @@ -66,19 +67,33 @@ defmodule ExCubicIngestion.ProcessIncoming do |> CubicTable.filter_to_existing_prefixes() for {table_prefix, table} <- table_prefixes do - incoming_bucket - |> S3Scan.list_objects_v2( - prefix: "#{incoming_prefix}#{table_prefix}", - lib_ex_aws: state.lib_ex_aws - ) - # filter s3 objects to only data objects with a size specified - |> Enum.filter(&Validators.valid_s3_object?(&1)) - |> Enum.map(fn object -> - %{object | key: String.replace_prefix(object[:key], incoming_prefix, "")} - end) - |> CubicLoad.insert_new_from_objects_with_table(table) + result = + incoming_bucket + |> S3Scan.list_objects_v2( + prefix: "#{incoming_prefix}#{table_prefix}", + lib_ex_aws: state.lib_ex_aws + ) + # filter s3 objects to only data objects with a size specified + |> Enum.filter(&Validators.valid_s3_object?(&1)) + |> Enum.map(fn object -> + %{object | key: String.replace_prefix(object[:key], incoming_prefix, "")} + end) + |> CubicLoad.insert_new_from_objects_with_table(table) + end + + case result do + {:ok, {_snapshot, [_ | _] = new_loads}} -> + Logger.info( + "parent=data_platform, process=run, table=#{table.name}, " <> + "s3_prefix=#{table.s3_prefix}, file_count=#{length(new_loads)}, " + ) + + _ -> + :ok end + result + :ok end diff --git a/py_cubic_ingestion/lib/py_cubic_ingestion/ingest_incoming.py b/py_cubic_ingestion/lib/py_cubic_ingestion/ingest_incoming.py index 8c16277..9ba9ff6 100644 --- a/py_cubic_ingestion/lib/py_cubic_ingestion/ingest_incoming.py +++ b/py_cubic_ingestion/lib/py_cubic_ingestion/ingest_incoming.py @@ -9,6 +9,7 @@ from py_cubic_ingestion import job_helpers from pyspark.context import SparkContext import boto3 +import logging import sys @@ -59,4 +60,13 @@ def run() -> None: # write out to springboard bucket using the same prefix as incoming job_helpers.write_parquet(updated_table_df, load.get("partition_columns", []), load["destination_path"]) + logging.info( + "parent=data_platform, process=ingest_incoming, source_table=%s, destination_table=%s, " + "destination_path=%s, source_s3_key=%s, row_count=%d, col_count=%d, " + "partition_count=%d", + load["source_table_name"], + load["destination_table_name"], + load["destination_path"], + load["source_s3_key"] + ) job.commit() From 078a3455be411f101386a0470fe7a16922d90189 Mon Sep 17 00:00:00 2001 From: Max Alexander Date: Thu, 5 Feb 2026 18:12:18 -0500 Subject: [PATCH 2/3] Add logging for files detected and files written --- .../ex_cubic_ingestion/process_incoming.ex | 38 +++++++++++++------ .../lib/py_cubic_ingestion/ingest_incoming.py | 10 +++++ 2 files changed, 37 insertions(+), 11 deletions(-) diff --git a/ex_cubic_ingestion/lib/ex_cubic_ingestion/process_incoming.ex b/ex_cubic_ingestion/lib/ex_cubic_ingestion/process_incoming.ex index 5ee2a73..731f73e 100644 --- a/ex_cubic_ingestion/lib/ex_cubic_ingestion/process_incoming.ex +++ b/ex_cubic_ingestion/lib/ex_cubic_ingestion/process_incoming.ex @@ -14,6 +14,7 @@ defmodule ExCubicIngestion.ProcessIncoming do alias ExCubicIngestion.Schema.CubicLoad alias ExCubicIngestion.Schema.CubicTable alias ExCubicIngestion.Validators + require Logger @wait_interval_ms 5_000 @@ -66,19 +67,34 @@ defmodule ExCubicIngestion.ProcessIncoming do |> CubicTable.filter_to_existing_prefixes() for {table_prefix, table} <- table_prefixes do - incoming_bucket - |> S3Scan.list_objects_v2( - prefix: "#{incoming_prefix}#{table_prefix}", - lib_ex_aws: state.lib_ex_aws - ) - # filter s3 objects to only data objects with a size specified - |> Enum.filter(&Validators.valid_s3_object?(&1)) - |> Enum.map(fn object -> - %{object | key: String.replace_prefix(object[:key], incoming_prefix, "")} - end) - |> CubicLoad.insert_new_from_objects_with_table(table) + result = + incoming_bucket + |> S3Scan.list_objects_v2( + prefix: "#{incoming_prefix}#{table_prefix}", + lib_ex_aws: state.lib_ex_aws + ) + # filter s3 objects to only data objects with a size specified + |> Enum.filter(&Validators.valid_s3_object?(&1)) + |> Enum.map(fn object -> + %{object | key: String.replace_prefix(object[:key], incoming_prefix, "")} + end) + |> CubicLoad.insert_new_from_objects_with_table(table) + + case result do + {:ok, {_snapshot, [_ | _] = new_loads}} -> + Logger.info( + "parent=data_platform, process=run, table=#{table.name}, " <> + "s3_prefix=#{table.s3_prefix}, file_count=#{length(new_loads)}, " + ) + + _ -> + :ok + end + + result end + :ok end diff --git a/py_cubic_ingestion/lib/py_cubic_ingestion/ingest_incoming.py b/py_cubic_ingestion/lib/py_cubic_ingestion/ingest_incoming.py index 8c16277..9ba9ff6 100644 --- a/py_cubic_ingestion/lib/py_cubic_ingestion/ingest_incoming.py +++ b/py_cubic_ingestion/lib/py_cubic_ingestion/ingest_incoming.py @@ -9,6 +9,7 @@ from py_cubic_ingestion import job_helpers from pyspark.context import SparkContext import boto3 +import logging import sys @@ -59,4 +60,13 @@ def run() -> None: # write out to springboard bucket using the same prefix as incoming job_helpers.write_parquet(updated_table_df, load.get("partition_columns", []), load["destination_path"]) + logging.info( + "parent=data_platform, process=ingest_incoming, source_table=%s, destination_table=%s, " + "destination_path=%s, source_s3_key=%s, row_count=%d, col_count=%d, " + "partition_count=%d", + load["source_table_name"], + load["destination_table_name"], + load["destination_path"], + load["source_s3_key"] + ) job.commit() From fa395033f071ebff8068798033533c0d426cb8f2 Mon Sep 17 00:00:00 2001 From: Max Alexander Date: Fri, 6 Feb 2026 17:04:23 -0500 Subject: [PATCH 3/3] Ran mix format check properly this time --- .../ex_cubic_ingestion/process_incoming.ex | 25 +++++++++---------- 1 file changed, 12 insertions(+), 13 deletions(-) diff --git a/ex_cubic_ingestion/lib/ex_cubic_ingestion/process_incoming.ex b/ex_cubic_ingestion/lib/ex_cubic_ingestion/process_incoming.ex index 731f73e..c685718 100644 --- a/ex_cubic_ingestion/lib/ex_cubic_ingestion/process_incoming.ex +++ b/ex_cubic_ingestion/lib/ex_cubic_ingestion/process_incoming.ex @@ -80,21 +80,20 @@ defmodule ExCubicIngestion.ProcessIncoming do end) |> CubicLoad.insert_new_from_objects_with_table(table) - case result do - {:ok, {_snapshot, [_ | _] = new_loads}} -> - Logger.info( - "parent=data_platform, process=run, table=#{table.name}, " <> - "s3_prefix=#{table.s3_prefix}, file_count=#{length(new_loads)}, " - ) - - _ -> - :ok - end - - result + case result do + {:ok, {_snapshot, [_ | _] = new_loads}} -> + Logger.info( + "parent=data_platform, process=run, table=#{table.name}, " <> + "s3_prefix=#{table.s3_prefix}, file_count=#{length(new_loads)}, " + ) + + _ -> + :ok + end + + result end - :ok end