Skip to content

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 image_path and epsg columns).

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
def __init__(self, spark: Any, config: SparkChippingConfig):
    """
    Args:
        spark: Active SparkSession.
        config: Chipping configuration.
    """
    self.spark = spark
    self.config = config

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.

Source code in src/geoembed/chipping/spark.py
def run(self) -> None:
    """
    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.
    """
    start_time = time.monotonic()
    quoted_input = quote_table_name(self.config.input_table)
    quoted_output = quote_table_name(self.config.output_table)

    source_df = self.spark.read.table(quoted_input)
    num_images = source_df.count()
    num_partitions = max(1, math.ceil(num_images / self.config.images_per_task))

    srid = self._resolve_srid(source_df)

    print(f"[SparkChippingPipeline] {num_images} images -> {num_partitions} partitions")
    print(f"[SparkChippingPipeline] Chip size: {self.config.chip_size}px")

    worker = _ChipMetadataWorker(chip_size=self.config.chip_size)
    result_df = source_df.repartition(num_partitions, "image_path").mapInPandas(
        worker, schema=OUTPUT_SCHEMA
    )

    from pyspark.sql.functions import expr

    result_df = result_df.withColumn(
        "geometry",
        expr(f"ST_GeomFromWKB(geometry_wkb, {srid})"),
    )

    (
        result_df.write.format("delta")
        .mode("overwrite")
        .option("overwriteSchema", "true")
        .saveAsTable(quoted_output)
    )

    self.spark.sql(f"ALTER TABLE {quoted_output} CLUSTER BY (chip_id)")

    elapsed = time.monotonic() - start_time
    row_count = self.spark.read.table(quoted_output).count()
    print(f"[SparkChippingPipeline] Complete: {row_count:,} chips in {elapsed:.1f}s")
    print(f"[SparkChippingPipeline] Output: {self.config.output_table}")