Distributed
Distributed computing: Spark utilities, HDFS operations and configuration.
Distributed functions package — lazy-loaded.
Contains Spark utilities, HDFS configuration, and HDFS operations. All submodules load on first attribute access via PEP 562 __getattr__.
- siege_utilities.distributed.sanitise_dataframe_column_names(df)[source]
Cleans dataframe column names by converting them to lowercase and replacing slashes/spaces with underscores.
- Parameters:
df (DataFrame) – Input Spark DataFrame.
- Returns:
Sanitised DataFrame.
- Return type:
DataFrame
- siege_utilities.distributed.tabulate_null_vs_not_null(df, column_name)[source]
Returns a dataframe showing the count of null and non-null values for a given column.
- Parameters:
df (DataFrame) – Input Spark DataFrame.
column_name (str) – Name of the column to analyze.
- Returns:
Resulting DataFrame with null vs non-null counts.
- Return type:
DataFrame
- siege_utilities.distributed.get_row_count(df)[source]
Returns the count of rows in the dataframe.
- Parameters:
df (DataFrame) – Input Spark DataFrame.
- Returns:
Row count.
- Return type:
- siege_utilities.distributed.repartition_and_cache(df, partitions=100)[source]
Repartitions and caches a dataframe.
- Parameters:
df (DataFrame) – Input Spark DataFrame.
partitions (int, optional) – Number of partitions. Default is 100.
- Returns:
Repartitioned and cached DataFrame.
- Return type:
DataFrame
- siege_utilities.distributed.register_temp_table(df, table_name)[source]
Registers a temporary view from a dataframe.
- Parameters:
df (DataFrame) – Input Spark DataFrame.
table_name (str) – Name for the temporary view.
- Return type:
None
- siege_utilities.distributed.move_column_to_front_of_dataframe(df, column_name)[source]
Reorder df so column_name is the leftmost column.
The Spark schema order is what most CSV / parquet readers surface first, and downstream consumers (notebook displays, exports) read left-to-right. Moving the join key / identifier to the front makes the rest of the pipeline more readable without changing semantics.
The original df is unchanged.
- Parameters:
df (None)
column_name (str)
- Return type:
None
- siege_utilities.distributed.write_df_to_parquet(df, path, mode='overwrite')[source]
Writes a DataFrame to a Parquet file.
- siege_utilities.distributed.read_parquet_to_df(spark, path)[source]
Reads a Parquet file into a Spark DataFrame.
- Parameters:
spark (SparkSession) – Active Spark session.
path (str) – Path to the Parquet file.
- Returns:
Loaded DataFrame.
- Return type:
DataFrame
- siege_utilities.distributed.flatten_json_column_and_join_back_to_df(df, json_column, prefix='json_column_', logger=None, drop_original=True, explode_arrays=False, flatten_level='shallow', verbose=False, sample_size=5, show_samples=True)[source]
Flattens a JSON column in a Spark DataFrame, extracting fields and adding them as columns. Has fallback mechanisms for corrupt JSON data.
- Parameters:
df (DataFrame) – The input Spark DataFrame.
json_column (str) – The name of the column containing JSON strings.
prefix (str, optional) – Prefix to add to the flattened column names. Defaults to “json_column_”.
logger (Optional[Any], optional) – Logger object for logging messages. Defaults to None.
drop_original (bool, optional) – Whether to drop the original JSON column after flattening. Defaults to True.
explode_arrays (bool, optional) – Whether to explode array columns. Defaults to False.
flatten_level (str, optional) – “shallow” or “deep” flattening. Defaults to “shallow”.
verbose (bool, optional) – Controls whether to log detailed messages. Defaults to False.
sample_size (int, optional) – Number of samples to check. Defaults to 5.
show_samples (bool, optional) – Whether to display sample data. Defaults to False.
- Returns:
The DataFrame with the JSON column flattened.
- Return type:
DataFrame
- Raises:
ValueError – If all JSON samples are corrupt and schema cannot be inferred.
RuntimeError – If Spark analysis fails and the fallback path also fails.
- siege_utilities.distributed.validate_geocode_data(df, lat_col_name, lon_col_name)[source]
Filters out rows with invalid geographic coordinates using string-based column names.
- Raises:
ValueError – If lat_col_name or lon_col_name is not in the DataFrame.
- Parameters:
- siege_utilities.distributed.mark_valid_geocode_data(df, lat_col_name, lon_col_name, output_col_name='is_valid')[source]
Adds a boolean flag column to the DataFrame indicating whether the geographic coordinates are valid.
A set of coordinates is considered valid if: - The latitude and longitude columns are not null. - The latitude is between -90 and 90. - The longitude is between -180 and 180.
Unlike filtering functions, this function preserves all rows in the DataFrame by simply marking each row with a True (valid) or False (invalid) value in the new output column.
- Parameters:
- Raises:
ValueError – If lat_col_name or lon_col_name is not in the DataFrame.
- Returns:
A new DataFrame with an additional column indicating geocode validity.
- Return type:
DataFrame
- siege_utilities.distributed.clean_and_reorder_bbox(df, bbox_col)[source]
Removes brackets from bounding box strings and reorders coordinates for Sedona.
- Assumes input is a comma separated list in the order:
min latitude, max latitude, min longitude, max longitude
Produces an array in the order: [min_lon, min_lat, max_lon, max_lat]
- siege_utilities.distributed.ensure_literal(value)[source]
Convert any value to a Spark literal (Column) unless it is already a Spark Column.
- Parameters:
value – Any value to be converted.
- Returns:
A pyspark.sql.Column containing the value (or its Spark literal), unless the value is already a Column.
- Return type:
None
- siege_utilities.distributed.reproject_geom_columns(df, geom_columns, source_srid, target_srid)[source]
Reprojects geometry columns using the three-argument version of ST_Transform: ST_Transform(geom, ‘source_srid’, ‘target_srid’)
Only reprojects if the current SRID is not equal to the target.
- Parameters:
- Raises:
ValueError – If source_srid or target_srid is not a valid EPSG identifier.
- Returns:
The DataFrame with each specified geometry column conditionally reprojected.
- Return type:
DataFrame
- siege_utilities.distributed.prepare_dataframe_for_export(df, logger_func=None)[source]
- Prepares a DataFrame for export (e.g., to CSV) by:
Converting binary columns to Base64-encoded strings.
Casting simple scalar fields (non-string, non-complex) to strings.
Dropping intermediate columns (e.g., ‘parsed_json’) if present.
Converting complex (StructType/ArrayType) columns to JSON strings.
Handling null values appropriately.
- Parameters:
df – Spark DataFrame to prepare
logger_func – Optional logging function (defaults to print)
- Returns:
The transformed DataFrame with all columns as strings or JSON strings.
- siege_utilities.distributed.prepare_summary_dataframe(data_tuples, column_names=None, logger_func=None)[source]
Helper function to create summary DataFrames with consistent string types. Prevents type merging errors by ensuring all values are strings.
- Parameters:
data_tuples – List of tuples with data
column_names – Column names for the DataFrame
logger_func – Optional logging function
- Raises:
RuntimeError – If no active Spark session is found.
- Returns:
Spark DataFrame with all string columns
- siege_utilities.distributed.export_pyspark_df_to_excel(df, file_name='output.xlsx', sheet_name='Sheet1')[source]
Converts a PySpark DataFrame to a Pandas DataFrame and exports it to an Excel file.
- siege_utilities.distributed.pivot_summary_table_for_bools(df, columns, spark)[source]
Generate a pivot table summary for given boolean flag columns in a DataFrame. The pivot table includes three metrics:
“Count”: Sum of rows where the flag is True.
“Percentage (%)”: Percentage relative to total records.
“Total”: The total number of records (repeated for each column).
All numeric values are converted to float to ensure a consistent type.
- Parameters:
df (DataFrame) – The source Spark DataFrame.
columns (list) – List of column names (assumed to be boolean flags) to summarize.
spark (SparkSession) – The active Spark session.
- Returns:
A Spark DataFrame representing the pivot table.
- Return type:
DataFrame
- siege_utilities.distributed.pivot_summary_with_metrics(df, group_col, pivot_col, spark)[source]
Generate a pivot summary for a categorical column against one or more grouping columns, including rows for “Count”, “Percentage (%)”, and “Total” for each group.
- Parameters:
df (DataFrame) – The source Spark DataFrame.
group_col (str or list) – The column name (or list of column names) used for grouping. For example, “geocode_granularity” or [“state”, “region”].
pivot_col (str) – The categorical column to pivot on (e.g., “final_geocode_choice”).
spark (SparkSession) – The active Spark session.
- Returns:
- A Spark DataFrame in which each original group appears as three rows:
one for the counts, one for the percentages, and one for the total count. The non-grouping columns represent each distinct pivot column value.
- Return type:
DataFrame
- siege_utilities.distributed.export_prepared_df_as_csv_to_path_using_delimiter(df, write_path, delimiter=',')[source]
Exports DataFrame with necessary transformations to ensure Spark compatibility.
- Parameters:
- Return type:
None
Applies prepare_dataframe_for_export() to prevent Spark export issues.
- siege_utilities.distributed.print_debug_table(spark_df, title)[source]
Helper function to convert a Spark DataFrame into a Pandas DataFrame, format it using tabulate, and print the result with a title.
- siege_utilities.distributed.compute_walkability(distance)[source]
Bucket a distance-in-meters into a walkability grade label.
Used to classify how far one place is from another in pedestrian terms — the thresholds come from urban-planning literature: Trivial (<100m), Tolerable, Moderate, Borderline, Outside (>500m). Returns a
{"grade": ..., "label": ...}dict, orNonewhen distance isNone(so the caller’s.withColumn(... apply udf)produces a null instead of crashing).The actual thresholds live in
walkability_configabove for audit / per-deployment tweaking.
- siege_utilities.distributed.validate_geometry(df, geom_col, step_name)[source]
Validates a single geometry column.
Parameters: - df (DataFrame): Spark DataFrame containing geometry data. - geom_col (str): Name of the geometry column to check. - step_name (str): Label for the debug output.
- siege_utilities.distributed.backup_full_dataframe(df, step_name)[source]
Persist df to
DEBUG_SUBDIRECTORY/{step_name}_full_persisted.A debug-mode helper: long Spark pipelines occasionally need an out-of-band snapshot of an intermediate frame (the canonical
.cache()lives in memory and dies with the job). This writes the snapshot to the configured debug directory in the project’s standard output format so it can be inspected after the run.No return value — purely a side-effect (write + log line).
- Parameters:
step_name (str)
- Return type:
None
- siege_utilities.distributed.atomic_write_with_staging(df, final_destination, staging_directory, file_format='csv', delimiter=',', header=True, mode='overwrite')[source]
Performs atomic write operations using a staging directory to prevent partial/corrupted files.
- siege_utilities.distributed.create_unique_staging_directory(base_path, operation_name='operation')[source]
Creates a unique staging directory for atomic operations.
- siege_utilities.distributed.py_round(number, ndigits=None)
Round a number to a given precision in decimal digits.
The return value is an integer if ndigits is omitted or None. Otherwise the return value has the same type as the number. ndigits may be negative.
- class siege_utilities.distributed.HDFSConfig[source]
Bases:
objectConfiguration class for HDFS and Spark operations.
Supports local, standalone cluster, and YARN deployments.
- __init__(data_path=None, hdfs_base_directory='/data/', cache_directory=None, app_name='SparkDistributedProcessing', spark_log_level='WARN', master='local[*]', deploy_mode='client', yarn_queue=None, enable_sedona=True, sedona_global_index_type='rtree', sedona_join_broadcast_threshold=10485760, num_executors=None, executor_cores=None, executor_memory='2g', driver_memory='2g', driver_cores=1, network_timeout='800s', heartbeat_interval='60s', hdfs_timeout=10, hdfs_copy_timeout=300, force_sync=False, log_info_func=None, log_error_func=None, hash_func=None, quick_signature_func=None)
- Parameters:
data_path (str | None)
hdfs_base_directory (str)
cache_directory (str | None)
app_name (str)
spark_log_level (str)
master (str)
deploy_mode (str)
yarn_queue (str | None)
enable_sedona (bool)
sedona_global_index_type (str)
sedona_join_broadcast_threshold (int)
num_executors (int | None)
executor_cores (int | None)
executor_memory (str)
driver_memory (str)
driver_cores (int)
network_timeout (str)
heartbeat_interval (str)
hdfs_timeout (int)
hdfs_copy_timeout (int)
force_sync (bool)
- Return type:
None
- siege_utilities.distributed.create_hdfs_config(**kwargs)[source]
Factory function to create HDFS configuration.
- Parameters:
**kwargs – Forwarded to HDFSConfig(). See HDFSConfig for accepted fields (e.g. data_path, master, executor_memory, enable_sedona, etc.).
- Return type:
- siege_utilities.distributed.create_local_config(data_path, **kwargs)[source]
Create config optimized for local development.
- Parameters:
data_path (str) – Path to data directory.
**kwargs – Override any HDFSConfig default. Merged on top of local-mode defaults (num_executors=2, executor_cores=1, executor_memory=’1g’, enable_sedona=False). See HDFSConfig for accepted fields.
- Return type:
- siege_utilities.distributed.create_cluster_config(data_path, **kwargs)[source]
Create config optimized for cluster deployment.
- Parameters:
data_path (str) – Path to data directory.
**kwargs – Override any HDFSConfig default. Merged on top of cluster defaults (num_executors=8, executor_cores=4, executor_memory=’4g’, enable_sedona=True). See HDFSConfig for accepted fields.
- Return type:
- siege_utilities.distributed.create_geocoding_config(data_path, **kwargs)[source]
Create config optimized for geocoding workloads.
- Parameters:
data_path (str) – Path to data directory.
**kwargs – Override any HDFSConfig default. Merged on top of geocoding defaults (app_name=’GeocodingPipeline’, num_executors=4, executor_memory=’2g’, enable_sedona=True, network_timeout=’1200s’). See HDFSConfig for accepted fields.
- Return type:
- siege_utilities.distributed.create_yarn_config(data_path, **kwargs)[source]
Create config optimized for YARN cluster deployment.
This configuration is designed for production YARN clusters with: - YARN as the resource manager - Client deploy mode (driver runs locally) - Higher resource allocation for executors - Sedona enabled for spatial operations
- Parameters:
data_path (str) – Path to data (local or HDFS)
**kwargs – Override any default settings
- Returns:
HDFSConfig configured for YARN deployment
- Return type:
Example
>>> config = create_yarn_config('/data/census', yarn_queue='analytics') >>> spark, path, _ = setup_distributed_environment(config)
- siege_utilities.distributed.create_census_analysis_config(data_path, **kwargs)[source]
Create config optimized for Census longitudinal analysis.
Designed for processing large Census datasets with: - Sedona for spatial operations (boundary crosswalks) - Optimized partitioning for Census tract-level data - Memory settings for time-series joins
- Parameters:
data_path (str) – Path to Census data
**kwargs – Override any default settings
- Returns:
HDFSConfig optimized for Census analysis
- Return type:
Example
>>> config = create_census_analysis_config('/data/census/acs') >>> spark, path, _ = setup_distributed_environment(config)
- class siege_utilities.distributed.AbstractHDFSOperations[source]
Bases:
objectAbstract HDFS Operations class that can be configured for any project
- create_spark_session()[source]
Create Spark session using configuration.
Supports local, standalone cluster, and YARN deployments based on the master URL in the config.
- sync_directory_to_hdfs(local_path=None, hdfs_subdir='inputs')[source]
Sync local directory/file to HDFS with proper verification.
- Raises:
ValueError – If no data path is provided.
FileNotFoundError – If the local path does not exist.
RuntimeError – If HDFS is not accessible or sync verification fails.
subprocess.CalledProcessError – If an HDFS command fails.
subprocess.TimeoutExpired – If an HDFS command times out.
- Parameters:
- Return type:
- setup_distributed_environment(data_path=None, dependency_paths=None)[source]
Main setup function with proper verification.
- Returns:
Tuple of (spark_session, data_url, None).
- Raises:
ValueError – If no data path is provided.
FileNotFoundError – If the data path does not exist.
ImportError – If PySpark is not available.
RuntimeError – If HDFS sync or Spark session creation fails.
- Parameters:
- siege_utilities.distributed.setup_distributed_environment(config, data_path=None, dependency_paths=None)[source]
Convenience function to set up distributed environment
- siege_utilities.distributed.create_hdfs_operations(config)[source]
Factory function to create HDFS operations instance
Submodules
- siege_utilities.distributed.spark_utils.sanitise_dataframe_column_names(df)[source]
Cleans dataframe column names by converting them to lowercase and replacing slashes/spaces with underscores.
- Parameters:
df (DataFrame) – Input Spark DataFrame.
- Returns:
Sanitised DataFrame.
- Return type:
DataFrame
- siege_utilities.distributed.spark_utils.tabulate_null_vs_not_null(df, column_name)[source]
Returns a dataframe showing the count of null and non-null values for a given column.
- Parameters:
df (DataFrame) – Input Spark DataFrame.
column_name (str) – Name of the column to analyze.
- Returns:
Resulting DataFrame with null vs non-null counts.
- Return type:
DataFrame
- siege_utilities.distributed.spark_utils.get_row_count(df)[source]
Returns the count of rows in the dataframe.
- Parameters:
df (DataFrame) – Input Spark DataFrame.
- Returns:
Row count.
- Return type:
- siege_utilities.distributed.spark_utils.repartition_and_cache(df, partitions=100)[source]
Repartitions and caches a dataframe.
- Parameters:
df (DataFrame) – Input Spark DataFrame.
partitions (int, optional) – Number of partitions. Default is 100.
- Returns:
Repartitioned and cached DataFrame.
- Return type:
DataFrame
- siege_utilities.distributed.spark_utils.register_temp_table(df, table_name)[source]
Registers a temporary view from a dataframe.
- Parameters:
df (DataFrame) – Input Spark DataFrame.
table_name (str) – Name for the temporary view.
- Return type:
None
- siege_utilities.distributed.spark_utils.move_column_to_front_of_dataframe(df, column_name)[source]
Reorder df so column_name is the leftmost column.
The Spark schema order is what most CSV / parquet readers surface first, and downstream consumers (notebook displays, exports) read left-to-right. Moving the join key / identifier to the front makes the rest of the pipeline more readable without changing semantics.
The original df is unchanged.
- Parameters:
df (None)
column_name (str)
- Return type:
None
- siege_utilities.distributed.spark_utils.write_df_to_parquet(df, path, mode='overwrite')[source]
Writes a DataFrame to a Parquet file.
- siege_utilities.distributed.spark_utils.read_parquet_to_df(spark, path)[source]
Reads a Parquet file into a Spark DataFrame.
- Parameters:
spark (SparkSession) – Active Spark session.
path (str) – Path to the Parquet file.
- Returns:
Loaded DataFrame.
- Return type:
DataFrame
- siege_utilities.distributed.spark_utils.flatten_json_column_and_join_back_to_df(df, json_column, prefix='json_column_', logger=None, drop_original=True, explode_arrays=False, flatten_level='shallow', verbose=False, sample_size=5, show_samples=True)[source]
Flattens a JSON column in a Spark DataFrame, extracting fields and adding them as columns. Has fallback mechanisms for corrupt JSON data.
- Parameters:
df (DataFrame) – The input Spark DataFrame.
json_column (str) – The name of the column containing JSON strings.
prefix (str, optional) – Prefix to add to the flattened column names. Defaults to “json_column_”.
logger (Optional[Any], optional) – Logger object for logging messages. Defaults to None.
drop_original (bool, optional) – Whether to drop the original JSON column after flattening. Defaults to True.
explode_arrays (bool, optional) – Whether to explode array columns. Defaults to False.
flatten_level (str, optional) – “shallow” or “deep” flattening. Defaults to “shallow”.
verbose (bool, optional) – Controls whether to log detailed messages. Defaults to False.
sample_size (int, optional) – Number of samples to check. Defaults to 5.
show_samples (bool, optional) – Whether to display sample data. Defaults to False.
- Returns:
The DataFrame with the JSON column flattened.
- Return type:
DataFrame
- Raises:
ValueError – If all JSON samples are corrupt and schema cannot be inferred.
RuntimeError – If Spark analysis fails and the fallback path also fails.
- siege_utilities.distributed.spark_utils.validate_geocode_data(df, lat_col_name, lon_col_name)[source]
Filters out rows with invalid geographic coordinates using string-based column names.
- Raises:
ValueError – If lat_col_name or lon_col_name is not in the DataFrame.
- Parameters:
- siege_utilities.distributed.spark_utils.mark_valid_geocode_data(df, lat_col_name, lon_col_name, output_col_name='is_valid')[source]
Adds a boolean flag column to the DataFrame indicating whether the geographic coordinates are valid.
A set of coordinates is considered valid if: - The latitude and longitude columns are not null. - The latitude is between -90 and 90. - The longitude is between -180 and 180.
Unlike filtering functions, this function preserves all rows in the DataFrame by simply marking each row with a True (valid) or False (invalid) value in the new output column.
- Parameters:
- Raises:
ValueError – If lat_col_name or lon_col_name is not in the DataFrame.
- Returns:
A new DataFrame with an additional column indicating geocode validity.
- Return type:
DataFrame
- siege_utilities.distributed.spark_utils.clean_and_reorder_bbox(df, bbox_col)[source]
Removes brackets from bounding box strings and reorders coordinates for Sedona.
- Assumes input is a comma separated list in the order:
min latitude, max latitude, min longitude, max longitude
Produces an array in the order: [min_lon, min_lat, max_lon, max_lat]
- siege_utilities.distributed.spark_utils.ensure_literal(value)[source]
Convert any value to a Spark literal (Column) unless it is already a Spark Column.
- Parameters:
value – Any value to be converted.
- Returns:
A pyspark.sql.Column containing the value (or its Spark literal), unless the value is already a Column.
- Return type:
None
- siege_utilities.distributed.spark_utils.reproject_geom_columns(df, geom_columns, source_srid, target_srid)[source]
Reprojects geometry columns using the three-argument version of ST_Transform: ST_Transform(geom, ‘source_srid’, ‘target_srid’)
Only reprojects if the current SRID is not equal to the target.
- Parameters:
- Raises:
ValueError – If source_srid or target_srid is not a valid EPSG identifier.
- Returns:
The DataFrame with each specified geometry column conditionally reprojected.
- Return type:
DataFrame
- siege_utilities.distributed.spark_utils.prepare_dataframe_for_export(df, logger_func=None)[source]
- Prepares a DataFrame for export (e.g., to CSV) by:
Converting binary columns to Base64-encoded strings.
Casting simple scalar fields (non-string, non-complex) to strings.
Dropping intermediate columns (e.g., ‘parsed_json’) if present.
Converting complex (StructType/ArrayType) columns to JSON strings.
Handling null values appropriately.
- Parameters:
df – Spark DataFrame to prepare
logger_func – Optional logging function (defaults to print)
- Returns:
The transformed DataFrame with all columns as strings or JSON strings.
- siege_utilities.distributed.spark_utils.prepare_summary_dataframe(data_tuples, column_names=None, logger_func=None)[source]
Helper function to create summary DataFrames with consistent string types. Prevents type merging errors by ensuring all values are strings.
- Parameters:
data_tuples – List of tuples with data
column_names – Column names for the DataFrame
logger_func – Optional logging function
- Raises:
RuntimeError – If no active Spark session is found.
- Returns:
Spark DataFrame with all string columns
- siege_utilities.distributed.spark_utils.export_pyspark_df_to_excel(df, file_name='output.xlsx', sheet_name='Sheet1')[source]
Converts a PySpark DataFrame to a Pandas DataFrame and exports it to an Excel file.
- siege_utilities.distributed.spark_utils.pivot_summary_table_for_bools(df, columns, spark)[source]
Generate a pivot table summary for given boolean flag columns in a DataFrame. The pivot table includes three metrics:
“Count”: Sum of rows where the flag is True.
“Percentage (%)”: Percentage relative to total records.
“Total”: The total number of records (repeated for each column).
All numeric values are converted to float to ensure a consistent type.
- Parameters:
df (DataFrame) – The source Spark DataFrame.
columns (list) – List of column names (assumed to be boolean flags) to summarize.
spark (SparkSession) – The active Spark session.
- Returns:
A Spark DataFrame representing the pivot table.
- Return type:
DataFrame
- siege_utilities.distributed.spark_utils.pivot_summary_with_metrics(df, group_col, pivot_col, spark)[source]
Generate a pivot summary for a categorical column against one or more grouping columns, including rows for “Count”, “Percentage (%)”, and “Total” for each group.
- Parameters:
df (DataFrame) – The source Spark DataFrame.
group_col (str or list) – The column name (or list of column names) used for grouping. For example, “geocode_granularity” or [“state”, “region”].
pivot_col (str) – The categorical column to pivot on (e.g., “final_geocode_choice”).
spark (SparkSession) – The active Spark session.
- Returns:
- A Spark DataFrame in which each original group appears as three rows:
one for the counts, one for the percentages, and one for the total count. The non-grouping columns represent each distinct pivot column value.
- Return type:
DataFrame
- siege_utilities.distributed.spark_utils.export_prepared_df_as_csv_to_path_using_delimiter(df, write_path, delimiter=',')[source]
Exports DataFrame with necessary transformations to ensure Spark compatibility.
- Parameters:
- Return type:
None
Applies prepare_dataframe_for_export() to prevent Spark export issues.
- siege_utilities.distributed.spark_utils.print_debug_table(spark_df, title)[source]
Helper function to convert a Spark DataFrame into a Pandas DataFrame, format it using tabulate, and print the result with a title.
- siege_utilities.distributed.spark_utils.compute_walkability(distance)[source]
Bucket a distance-in-meters into a walkability grade label.
Used to classify how far one place is from another in pedestrian terms — the thresholds come from urban-planning literature: Trivial (<100m), Tolerable, Moderate, Borderline, Outside (>500m). Returns a
{"grade": ..., "label": ...}dict, orNonewhen distance isNone(so the caller’s.withColumn(... apply udf)produces a null instead of crashing).The actual thresholds live in
walkability_configabove for audit / per-deployment tweaking.
- siege_utilities.distributed.spark_utils.validate_geometry(df, geom_col, step_name)[source]
Validates a single geometry column.
Parameters: - df (DataFrame): Spark DataFrame containing geometry data. - geom_col (str): Name of the geometry column to check. - step_name (str): Label for the debug output.
- siege_utilities.distributed.spark_utils.backup_full_dataframe(df, step_name)[source]
Persist df to
DEBUG_SUBDIRECTORY/{step_name}_full_persisted.A debug-mode helper: long Spark pipelines occasionally need an out-of-band snapshot of an intermediate frame (the canonical
.cache()lives in memory and dies with the job). This writes the snapshot to the configured debug directory in the project’s standard output format so it can be inspected after the run.No return value — purely a side-effect (write + log line).
- Parameters:
step_name (str)
- Return type:
None
- siege_utilities.distributed.spark_utils.atomic_write_with_staging(df, final_destination, staging_directory, file_format='csv', delimiter=',', header=True, mode='overwrite')[source]
Performs atomic write operations using a staging directory to prevent partial/corrupted files.
- siege_utilities.distributed.spark_utils.create_unique_staging_directory(base_path, operation_name='operation')[source]
Creates a unique staging directory for atomic operations.
Abstract HDFS Operations - Fully Configurable and Reusable Zero hard-coded project dependencies
- class siege_utilities.distributed.hdfs_operations.AbstractHDFSOperations[source]
Bases:
objectAbstract HDFS Operations class that can be configured for any project
- create_spark_session()[source]
Create Spark session using configuration.
Supports local, standalone cluster, and YARN deployments based on the master URL in the config.
- sync_directory_to_hdfs(local_path=None, hdfs_subdir='inputs')[source]
Sync local directory/file to HDFS with proper verification.
- Raises:
ValueError – If no data path is provided.
FileNotFoundError – If the local path does not exist.
RuntimeError – If HDFS is not accessible or sync verification fails.
subprocess.CalledProcessError – If an HDFS command fails.
subprocess.TimeoutExpired – If an HDFS command times out.
- Parameters:
- Return type:
- setup_distributed_environment(data_path=None, dependency_paths=None)[source]
Main setup function with proper verification.
- Returns:
Tuple of (spark_session, data_url, None).
- Raises:
ValueError – If no data path is provided.
FileNotFoundError – If the data path does not exist.
ImportError – If PySpark is not available.
RuntimeError – If HDFS sync or Spark session creation fails.
- Parameters:
- siege_utilities.distributed.hdfs_operations.create_hdfs_operations(config)[source]
Factory function to create HDFS operations instance
- siege_utilities.distributed.hdfs_operations.setup_distributed_environment(config, data_path=None, dependency_paths=None)[source]
Convenience function to set up distributed environment
Abstract HDFS Configuration Configurable settings for HDFS operations - no hard-coded project dependencies
- class siege_utilities.distributed.hdfs_config.HDFSConfig[source]
Bases:
objectConfiguration class for HDFS and Spark operations.
Supports local, standalone cluster, and YARN deployments.
- __init__(data_path=None, hdfs_base_directory='/data/', cache_directory=None, app_name='SparkDistributedProcessing', spark_log_level='WARN', master='local[*]', deploy_mode='client', yarn_queue=None, enable_sedona=True, sedona_global_index_type='rtree', sedona_join_broadcast_threshold=10485760, num_executors=None, executor_cores=None, executor_memory='2g', driver_memory='2g', driver_cores=1, network_timeout='800s', heartbeat_interval='60s', hdfs_timeout=10, hdfs_copy_timeout=300, force_sync=False, log_info_func=None, log_error_func=None, hash_func=None, quick_signature_func=None)
- Parameters:
data_path (str | None)
hdfs_base_directory (str)
cache_directory (str | None)
app_name (str)
spark_log_level (str)
master (str)
deploy_mode (str)
yarn_queue (str | None)
enable_sedona (bool)
sedona_global_index_type (str)
sedona_join_broadcast_threshold (int)
num_executors (int | None)
executor_cores (int | None)
executor_memory (str)
driver_memory (str)
driver_cores (int)
network_timeout (str)
heartbeat_interval (str)
hdfs_timeout (int)
hdfs_copy_timeout (int)
force_sync (bool)
- Return type:
None
- siege_utilities.distributed.hdfs_config.create_census_analysis_config(data_path, **kwargs)[source]
Create config optimized for Census longitudinal analysis.
Designed for processing large Census datasets with: - Sedona for spatial operations (boundary crosswalks) - Optimized partitioning for Census tract-level data - Memory settings for time-series joins
- Parameters:
data_path (str) – Path to Census data
**kwargs – Override any default settings
- Returns:
HDFSConfig optimized for Census analysis
- Return type:
Example
>>> config = create_census_analysis_config('/data/census/acs') >>> spark, path, _ = setup_distributed_environment(config)
- siege_utilities.distributed.hdfs_config.create_cluster_config(data_path, **kwargs)[source]
Create config optimized for cluster deployment.
- Parameters:
data_path (str) – Path to data directory.
**kwargs – Override any HDFSConfig default. Merged on top of cluster defaults (num_executors=8, executor_cores=4, executor_memory=’4g’, enable_sedona=True). See HDFSConfig for accepted fields.
- Return type:
- siege_utilities.distributed.hdfs_config.create_geocoding_config(data_path, **kwargs)[source]
Create config optimized for geocoding workloads.
- Parameters:
data_path (str) – Path to data directory.
**kwargs – Override any HDFSConfig default. Merged on top of geocoding defaults (app_name=’GeocodingPipeline’, num_executors=4, executor_memory=’2g’, enable_sedona=True, network_timeout=’1200s’). See HDFSConfig for accepted fields.
- Return type:
- siege_utilities.distributed.hdfs_config.create_hdfs_config(**kwargs)[source]
Factory function to create HDFS configuration.
- Parameters:
**kwargs – Forwarded to HDFSConfig(). See HDFSConfig for accepted fields (e.g. data_path, master, executor_memory, enable_sedona, etc.).
- Return type:
- siege_utilities.distributed.hdfs_config.create_local_config(data_path, **kwargs)[source]
Create config optimized for local development.
- Parameters:
data_path (str) – Path to data directory.
**kwargs – Override any HDFSConfig default. Merged on top of local-mode defaults (num_executors=2, executor_cores=1, executor_memory=’1g’, enable_sedona=False). See HDFSConfig for accepted fields.
- Return type:
- siege_utilities.distributed.hdfs_config.create_yarn_config(data_path, **kwargs)[source]
Create config optimized for YARN cluster deployment.
This configuration is designed for production YARN clusters with: - YARN as the resource manager - Client deploy mode (driver runs locally) - Higher resource allocation for executors - Sedona enabled for spatial operations
- Parameters:
data_path (str) – Path to data (local or HDFS)
**kwargs – Override any default settings
- Returns:
HDFSConfig configured for YARN deployment
- Return type:
Example
>>> config = create_yarn_config('/data/census', yarn_queue='analytics') >>> spark, path, _ = setup_distributed_environment(config)