Spark Chipping Pipeline¶
geoembed.chipping.spark
¶
Spark-based distributed chipping pipeline for Databricks.
Reads a metadata table, distributes chip geometry calculation across workers, writes results to a Delta table with Databricks native geometry for spatial SQL.
SparkChippingConfig(input_table, output_table, chip_size=224, images_per_task=10)
dataclass
¶
Configuration for the Spark chipping pipeline.
Attributes:
| Name | Type | Description |
|---|---|---|
input_table |
str
|
Source metadata table (must have |
output_table |
str
|
Delta table to write chip metadata to. |
chip_size |
int
|
Square chip dimension in pixels (default 224 for DOFA). |
images_per_task |
int
|
Approximate images per Spark task (controls partition count). |
SparkChippingPipeline(spark, config)
¶
Distributed chipping pipeline for Databricks.
Reads a metadata table, distributes chip geometry calculation across Spark workers using mapInPandas, adds a Databricks native geometry column, and writes results to a Delta table with liquid clustering on chip_id.
Usage
from geoembed.chipping.spark import SparkChippingPipeline, SparkChippingConfig
config = SparkChippingConfig( input_table="catalog.schema.metadata", output_table="catalog.schema.chip_metadata", chip_size=224, ) pipeline = SparkChippingPipeline(spark, config) pipeline.run()
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
spark
|
Any
|
Active SparkSession. |
required |
config
|
SparkChippingConfig
|
Chipping configuration. |
required |
Source code in src/geoembed/chipping/spark.py
run()
¶
Execute the chipping pipeline.
Reads source metadata, distributes chipping, adds native geometry (SRID derived from the metadata table's epsg column), writes to Delta, and applies liquid clustering.