Rasterflow, Earth Intelligence & inference engine now in public preview Learn More

What Is the Partitioning Tax? Tiles, Seams, and Reruns

Authors

The partitioning tax is the cost an analyst pays when spatial data outgrows one machine. The workflow gets split into tiles, often one per county, and the analyst inherits a loop: process each tile, repair the edges, stitch the outputs, and rerun when a parameter changes. Wherobots named the pattern in its Wherobots for QGIS announcement. This guide shows where the tax comes from, how much it costs on real Colorado data, and how a distributed spatial engine runs the same analysis in one query over the full extent.

Key takeaways

  • The partitioning tax has three costs: engineering time, edge effects at tile seams, and reruns of every tile.
  • Tile seams create edge effects, the boundary bias spatial statistics has corrected since the 1950s. A per-row filter or count has none. Edge effects hit anything that looks past a feature's own location: joins, nearest-neighbor and distance queries, buffers, focal raster operations, and spatial statistics.
  • In Colorado, a per-county loop gave 188 of 2,139 schools the wrong nearest hospital or none. One ST_KNN query over the full extent gave every school its true nearest hospital.
  • An overlap buffer fixes the seams by copying data. A 50 km halo removed every error and stored each hospital an average of 8.4 times.
  • WherobotsDB partitions space, indexes each partition, copies edge features, and removes duplicates inside the engine, for vector and raster data. A larger extent calls for a larger runtime and the same SQL.

What is the partitioning tax?

The partitioning tax is the overhead of splitting a spatial analysis into pieces that each fit one machine, and of repairing what the split breaks. Desktop GIS and single-node databases process data on one machine, so data that outgrows its memory gets split into tiles, merged, and debugged at the edges.

The tax has three costs:

  • Engineering time. Someone writes and maintains the loop: the tile grid, the per-tile job, the merge step, and the bookkeeping that tracks which tiles finished.
  • Correctness at the seams. Features that sit on a tile edge, or whose neighbors sit across it, come out dropped, duplicated, or wrong. These are edge effects, and they are silent: every tile succeeds, and the output has the right shape.
  • Reruns. A new buffer distance, a new search radius, or a new input layer means running every tile again, and re-tuning the seam handling for the new parameter.

Testing at city scale hides the second cost: a city sits inside one tile, so it shows no seams.

A query that handles each row on its own, such as a filter, a count, or an aggregate by region, tiles cleanly: each row lands in one tile. The tax comes due when an answer needs data from across a tile edge: usually a second dataset, or, for focal raster operations, neighboring cells. A join needs both sides of each matching pair in one tile, and pairs can span the edge. Apache Sedona's single-node SpatialBench results show the same split on one machine. DuckDB and SedonaDB "achieve similar low-latency performance" on spatial filters and basic operations at SF 1 and SF 10. On complex joins, "DuckDB handles some join queries well but encounters scaling issues in certain cases," and DuckDB and GeoPandas "currently lack built-in KNN join support."

Edge effects: the error behind the tax

An edge effect, also called a boundary effect, is the bias that appears when an analysis treats the edge of its data as the edge of the world. A feature near the boundary has part of its neighborhood outside the data, so anything computed from it comes out biased: nearest neighbors too far, counts too low, focal windows missing cells. A tile loop adds a boundary at every tile edge, so every seam error in this article is an edge effect.

The term has two meanings. In ecology, Aldo Leopold used "edge effect" in Game Management (1933) for the greater variety of wildlife where two habitats meet. In spatial statistics it names the boundary bias, and its corrections came from point-pattern analysis:

  • Nearest-neighbor distances. Philip Clark and Francis Evans's nearest-neighbor index (Ecology, 1954) compares observed nearest-neighbor distances with those of a random pattern. Near the boundary, a point's true nearest neighbor can lie outside the study area, so distances run long.
  • Weighting. Brian Ripley's K function, in Modelling spatial patterns (Journal of the Royal Statistical Society, Series B, 1977), weights each pair of points by the inverse of the share of its circle's circumference inside the study area, so a pair whose circle is half outside counts twice.
  • Buffer zones and wrapping. Daniel Griffith and Carl Amrhein's evaluation of correction techniques for boundary effects (Geographical Analysis, 1983) reviewed the traditional options: ignore the effect, map the study area onto a torus so opposite edges touch, build an empirical or artificial buffer zone, extrapolate into the buffer, or apply a correction factor.

The halo in a tile loop is the empirical buffer zone: read real data past the edge, report only the features inside. Weighting cannot replace the actual neighbor a join needs, and a torus suits simulations only. The buffer zone carries the cost the walkthrough measures: it has to be as wide as the farthest interaction, and for nearest-neighbor search that width depends on how sparse the data is. One query over the full extent leaves only the data's true boundary, such as a state line, which the walkthrough covers by reading hospitals in neighboring states.

How a tile loop breaks

Waldo Tobler's first law of geography (Economic Geography, 1970) holds that near things are more related than distant things, and a tile edge cuts through exactly those near relationships. Each break below is an edge effect.

Four-panel schematic of two tiles split by a dashed yellow edge. Panel 1: a school in tile A linked by a dashed red line to a far hospital in tile A and by a green line to a closer hospital in tile B. Panel 2: a polygon straddling the edge, labeled copy in A and copy in B. Panel 3: a halo band around the edge and two overlapping boxes whose intersection has a reference point marked in tile A. Panel 4: a 3 by 3 raster window on the last column of tile A that needs halo cells from tile B
Four ways a tile edge breaks a spatial analysis, and the halo plus reference point fix. Schematic.

Partitioning tax examples

Spatial joins and point in polygon

A spatial join matches rows by location. Assign polygons to tiles with ST_Intersects and a polygon that crosses the edge lands in two tiles and is counted twice. Keep only features that fall entirely inside a tile and it lands in neither. Roads, rivers, parcels, and building footprints often cross edges.

Distance and nearest-neighbor queries

A k-nearest-neighbor (KNN) query returns the closest features to each row. Inside a tile, the search covers only that tile's candidates, so a school near a county line gets the nearest hospital in its own county even when a closer one sits across the line. A school in a county with no hospital gets no answer. Distance joins such as ST_DWithin fail the same way across the edge.

Buffers, dissolves, and aggregations

A buffer drawn in one tile stops at the tile edge unless the neighbor's features are present. A dissolve that merges touching polygons leaves a seam line wherever the merged shape crossed tiles. Aggregations take on the tiles as zones: a count per county or per grid cell depends on where the boundaries fall, which is the modifiable areal unit problem. Stan Openshaw described it in The Modifiable Areal Unit Problem (CATMOG 38, 1984), and A. Stewart Fotheringham and David Wong showed that it makes multivariate statistical results unpredictable (Environment and Planning A, 1991). H3 cells give every zone the same shape, but counts still shift with cell size and grid placement, so H3 reduces the problem without removing it.

Focal raster operations

Slope, aspect, hillshade, and smoothing filters compute each cell from a window of its neighbors. Berthold Horn's slope and shading estimates use a 3 by 3 window (Hill shading and the reflectance map, Proceedings of the IEEE, 1981), so the last row and column of a tile have missing neighbors. A tiled raster needs halo cells from adjacent tiles, as the digital elevation model walkthrough adds with a one-cell border.

The usual fixes and what each one costs

Overlap buffers, or halos. The buffer-zone correction above: each tile reads data within a set distance past its edge. In raster work, these are halo or ghost cells.

Reference-point duplicate avoidance. A halo puts features in more than one tile, so results repeat. Jens-Peter Dittrich and Bernhard Seeger gave the standard rule in Data redundancy and duplicate detection in spatial join processing (ICDE 2000): report a pair only in the tile that contains a fixed corner of the intersection of the two bounding boxes. Each pair has one such corner, so each pair is reported once.

Deduplication. Without a reference-point rule, a final DISTINCT or GROUP BY on feature IDs removes copies. That needs a stable ID on every feature and a shuffle of the full output.

Tuning. Each fix has its own parameter (halo width, tile size, the dedupe key). A larger search radius needs a wider halo, which copies more rows, which changes how big a tile can be.

How distributed spatial engines remove the tax

A distributed spatial engine runs a careful tile loop's steps inside the query planner.

  • Spatial partitioning. The engine samples the data and builds partitions from it: a uniform grid, a quadtree (Finkel and Bentley, 1974), or a KDB-tree (Robinson, 1981), which splits at the median. WherobotsDB uses a KDB-tree for joins by default (sedona.join.gridtype = kdbtree in the WherobotsDB parameters).
  • Partition-level indexes. Each partition gets an R-tree (Guttman, 1984) or quadtree, so candidate pairs come from an index probe.
  • Replication with duplicate elimination. A geometry that crosses a partition edge is copied to each partition it touches, and duplicate results are removed before output, so the engine handles its own partition edges. Jia Yu, Zongsi Zhang, and Mohamed Sarwat describe these partitioners and per-partition indexes in Spatial data management in Apache Spark: the GeoSpark perspective and beyond (GeoInformatica, 2019).
  • Broadcast joins. When one side is small, such as 64 county polygons, it is copied to every worker with an index.

In WherobotsDB, EXPLAIN FORMATTED prints which of these the plan uses. The spatial join guide verified BroadcastIndexJoin, RangeJoin, DistanceJoin, and BroadcastQuerySideKNNJoin. The same planning applies to raster tiles joined with RS_Intersects. Wherobots runtimes range from General Purpose Micro to 4X-Large, plus Memory Optimized sizes; a larger extent takes a larger runtime and the same query.

The partitioning tax and AI coding assistants

The tax reaches AI coding assistants too. Their code comes from patterns in their training data (text, documents, databases, and the internet), and for a single-machine stack the scaling pattern is the loop, usually with no overlap buffer and no dedupe. Against an engine that handles the extent, the same request is one SQL statement. The Wherobots MCP server lets an assistant run it in Wherobots, and the QGIS plugin loads the result as a map layer.

Measuring the partitioning tax with Wherobots

WherobotsDB runs each analysis below as one query across the full extent, so it pays no partitioning tax. To measure the tax a tile loop pays, this walkthrough simulates the loop in SQL on one state and checks it against the full-extent answer. Colorado's 64 counties, from dense Front Range metros to sparse rural ones, show every type of seam. Three queries ran in a Wherobots SQL session on Overture Maps places, Overture buildings, and US Census TIGER/Line counties from Wherobots Open Data in the Havasu catalog. Every number below comes from the SQL printed with it.

Nearest hospital: a simulated tile loop vs. one query

A per-row count tiles cleanly. Counted county by county with the main query's point-in-polygon assignment, the schools add up to the statewide total:

-- One dataset, tiled by county: do the per-county counts add up?
WITH counties AS (
  SELECT GEOID AS county, geometry
  FROM wherobots_open_data.us_census.tiger_county
  WHERE STATEFP = '08'
),
schools AS (
  SELECT p.id, c.county
  FROM wherobots_open_data.overture_maps_foundation.places_place p
  JOIN counties c ON ST_Intersects(c.geometry, p.geometry)
  WHERE p.taxonomy.primary IN ('elementary_school', 'middle_school', 'high_school')
    AND p.bbox.xmin BETWEEN -109.1 AND -102.0 AND p.bbox.ymin BETWEEN 36.9 AND 41.1
),
per_county AS (SELECT county, COUNT(*) AS n FROM schools GROUP BY county)
SELECT (SELECT SUM(n) FROM per_county) AS sum_of_county_counts,
       (SELECT COUNT(DISTINCT id) FROM schools) AS distinct_schools,
       (SELECT COUNT(*) FROM per_county) AS counties_with_schools

The 62 county counts sum to 2,139, the number of distinct schools. The tax appears when a second dataset joins.

The question: which hospital is nearest to each school? Schools are Overture places with a primary category of elementary, middle, or high school, and hospitals are places with the category hospital, a label that includes some clinics and specialty centers. The tile_loop step simulates a county tile loop with no overlap buffer: each school searches only the hospitals in its own county, assigned by point in polygon. The full_extent step is the normal way to work in Wherobots: one ST_KNN join over every hospital in a box one degree wider than the state. Both sides use a coordinate reference system in meters as the nearest-neighbor documentation recommends for exact results.

-- Nearest hospital for every Colorado school, computed two ways
-- Distances in meters in a Lambert conformal conic projection centered on Colorado
WITH counties AS (
  SELECT GEOID AS county, NAME AS county_name, geometry
  FROM wherobots_open_data.us_census.tiger_county
  WHERE STATEFP = '08'
),
schools AS (
  -- Each school gets the county it falls in (point in polygon): its tile
  SELECT p.id, p.names.primary AS school, c.county, c.county_name,
         ST_X(p.geometry) AS lon, ST_Y(p.geometry) AS lat,
         ST_Transform(p.geometry, 'EPSG:4326', '+proj=lcc +lat_1=37.5 +lat_2=40.5 +lat_0=39 +lon_0=-105.5 +x_0=0 +y_0=0 +datum=NAD83 +units=m +no_defs') AS geom_m
  FROM wherobots_open_data.overture_maps_foundation.places_place p
  JOIN counties c ON ST_Intersects(c.geometry, p.geometry)
  WHERE p.taxonomy.primary IN ('elementary_school', 'middle_school', 'high_school')
    AND p.bbox.xmin BETWEEN -109.1 AND -102.0 AND p.bbox.ymin BETWEEN 36.9 AND 41.1
),
hospitals AS (
  -- Hospitals in a box one degree wider than the schools; county is null outside Colorado
  SELECT p.id, p.names.primary AS hospital, c.county,
         ST_X(p.geometry) AS lon, ST_Y(p.geometry) AS lat,
         ST_Transform(p.geometry, 'EPSG:4326', '+proj=lcc +lat_1=37.5 +lat_2=40.5 +lat_0=39 +lon_0=-105.5 +x_0=0 +y_0=0 +datum=NAD83 +units=m +no_defs') AS geom_m
  FROM wherobots_open_data.overture_maps_foundation.places_place p
  LEFT JOIN counties c ON ST_Intersects(c.geometry, p.geometry)
  WHERE p.taxonomy.primary = 'hospital'
    AND p.bbox.xmin BETWEEN -110.1 AND -101.0 AND p.bbox.ymin BETWEEN 35.9 AND 42.1
),
tile_loop AS (
  -- Simulated tile loop: each county searches only its own hospitals, with no overlap buffer
  SELECT id, hospital_id, hospital, distance_m, h_lon, h_lat FROM (
    SELECT s.id, h.id AS hospital_id, h.hospital, h.lon AS h_lon, h.lat AS h_lat,
           ST_Distance(s.geom_m, h.geom_m) AS distance_m,
           ROW_NUMBER() OVER (PARTITION BY s.id ORDER BY ST_Distance(s.geom_m, h.geom_m), h.id) AS rn
    FROM schools s JOIN hospitals h ON s.county = h.county
  ) WHERE rn = 1
),
full_extent AS (
  -- One pass over the full extent: k-nearest-neighbor join
  SELECT s.id, h.id AS hospital_id, h.hospital, h.county AS hospital_county,
         h.lon AS h_lon, h.lat AS h_lat, ST_Distance(s.geom_m, h.geom_m) AS distance_m
  FROM schools s JOIN hospitals h ON ST_KNN(s.geom_m, h.geom_m, 1, false)
),
cmp AS (
  SELECT s.id, s.school, s.county, s.county_name, s.lon, s.lat,
         f.hospital AS full_hospital, f.h_lon AS full_lon, f.h_lat AS full_lat,
         f.distance_m AS full_m, f.hospital_county,
         t.hospital AS tile_hospital, t.h_lon AS tile_lon, t.h_lat AS tile_lat, t.distance_m AS tile_m,
         CASE WHEN t.id IS NULL THEN 'no_hospital_in_county'
              WHEN t.distance_m > f.distance_m + 1 THEN 'wrong_hospital'
              ELSE 'same' END AS outcome
  FROM schools s
  JOIN full_extent f ON s.id = f.id
  LEFT JOIN tile_loop t ON s.id = t.id
)
SELECT *,
  COUNT(*) OVER () AS schools,
  SUM(CASE WHEN outcome = 'no_hospital_in_county' THEN 1 ELSE 0 END) OVER () AS n_dropped,
  SUM(CASE WHEN outcome = 'wrong_hospital' THEN 1 ELSE 0 END) OVER () AS n_wrong,
  SUM(CASE WHEN hospital_county IS NULL THEN 1 ELSE 0 END) OVER () AS n_nearest_out_of_state,
  SUM(CASE WHEN hospital_county IS NOT NULL AND hospital_county <> county THEN 1 ELSE 0 END) OVER () AS n_nearest_other_county,
  ROUND(AVG(CASE WHEN outcome = 'wrong_hospital' THEN tile_m - full_m END) OVER ()) AS mean_extra_m,
  ROUND(MAX(CASE WHEN outcome = 'wrong_hospital' THEN tile_m - full_m END) OVER ()) AS max_extra_m,
  ROUND(AVG(full_m) OVER ()) AS mean_full_m,
  SIZE(COLLECT_SET(county) OVER ()) AS counties_with_schools
FROM cmp
ORDER BY outcome, tile_m - full_m DESC

The query ran in 16.1 seconds and compared 2,139 schools in 62 counties. The per-county run gave 121 schools a farther hospital than the true nearest and gave 67 schools no hospital at all, because their county has no hospital record: 188 schools, 8.8%, with a wrong answer or none. These 188 schools are the edge effect at the county lines. For them, the true nearest hospital sits in another Colorado county for 181 and in another state for 7. Across the wrong answers the per-county hospital is 5,346 meters farther on average. The worst case is Prairie High School in Weld County: the per-county run returned a Weld County hospital 78.6 km away, and the true nearest, Fort Morgan Community Hospital Association in Morgan County, is 39.4 km away, a 39,168-meter error.

Map of the Denver metro with county outlines for Boulder, Adams, Denver, Jefferson, Arapahoe and Douglas. Grey dots mark schools with the same answer both ways. Red dots mark schools where the per-county run picked a farther hospital, each with a dashed red line to that hospital and a solid green line to the true nearest hospital across a county line. Side panel: 2,139 Colorado schools, 121 wrong hospital, 67 no hospital in their county, 39,168 m largest extra distance
Denver metro detail of the Colorado run. In the per-county simulation, 121 of 2,139 schools got a farther hospital (red) and 67 got none; the full-extent ST_KNN join found each true nearest (green).

The halo fix and its cost

The standard repair grows each county tile by a halo, so each school searches the hospitals in its county's grown tile. This query tries five widths at once and counts the errors left and the rows copied.

-- The usual fix: give every county tile an overlap buffer (halo) and count what it costs
WITH halos AS (
  SELECT * FROM VALUES (0), (5000), (10000), (25000), (50000) AS t(halo_m)
),
counties AS (
  SELECT GEOID AS county, geometry,
         ST_Transform(geometry, 'EPSG:4326', '+proj=lcc +lat_1=37.5 +lat_2=40.5 +lat_0=39 +lon_0=-105.5 +x_0=0 +y_0=0 +datum=NAD83 +units=m +no_defs') AS geom_m
  FROM wherobots_open_data.us_census.tiger_county
  WHERE STATEFP = '08'
),
tiles AS (
  -- One tile per county and halo width: the county polygon grown by halo_m meters
  SELECT /*+ BROADCAST(h) */ c.county, h.halo_m,
         CASE WHEN h.halo_m = 0 THEN c.geom_m ELSE ST_Buffer(c.geom_m, h.halo_m) END AS tile_m
  FROM counties c CROSS JOIN halos h
),
schools AS (
  SELECT p.id, c.county,
         ST_Transform(p.geometry, 'EPSG:4326', '+proj=lcc +lat_1=37.5 +lat_2=40.5 +lat_0=39 +lon_0=-105.5 +x_0=0 +y_0=0 +datum=NAD83 +units=m +no_defs') AS geom_m
  FROM wherobots_open_data.overture_maps_foundation.places_place p
  JOIN counties c ON ST_Intersects(c.geometry, p.geometry)
  WHERE p.taxonomy.primary IN ('elementary_school', 'middle_school', 'high_school')
    AND p.bbox.xmin BETWEEN -109.1 AND -102.0 AND p.bbox.ymin BETWEEN 36.9 AND 41.1
),
hospitals AS (
  SELECT p.id, ST_Transform(p.geometry, 'EPSG:4326', '+proj=lcc +lat_1=37.5 +lat_2=40.5 +lat_0=39 +lon_0=-105.5 +x_0=0 +y_0=0 +datum=NAD83 +units=m +no_defs') AS geom_m
  FROM wherobots_open_data.overture_maps_foundation.places_place p
  WHERE p.taxonomy.primary = 'hospital'
    AND p.bbox.xmin BETWEEN -110.1 AND -101.0 AND p.bbox.ymin BETWEEN 35.9 AND 42.1
),
tile_hospitals AS (
  -- Every hospital is copied into each tile it falls in
  SELECT t.county, t.halo_m, h.id, h.geom_m
  FROM tiles t JOIN hospitals h ON ST_Intersects(t.tile_m, h.geom_m)
),
tile_schools AS (
  -- What a loop reading every school in its grown tile would see; copies past the first are repeats to drop.
  -- (best below keeps each school in its home county, so it does not use these rows.)
  SELECT t.county, t.halo_m, s.id
  FROM tiles t JOIN schools s ON ST_Intersects(t.tile_m, s.geom_m)
),
best AS (
  -- Each school searches the hospitals of its own county's tile
  SELECT s.id, th.halo_m, MIN(ST_Distance(s.geom_m, th.geom_m)) AS tile_m
  FROM schools s JOIN tile_hospitals th ON s.county = th.county
  GROUP BY s.id, th.halo_m
),
truth AS (
  SELECT s.id, ST_Distance(s.geom_m, h.geom_m) AS full_m
  FROM schools s JOIN hospitals h ON ST_KNN(s.geom_m, h.geom_m, 1, false)
),
scored AS (
  SELECT sh.halo_m,
         COUNT(*) AS schools,
         SUM(CASE WHEN b.tile_m IS NULL THEN 1 ELSE 0 END) AS dropped,
         SUM(CASE WHEN b.tile_m > t.full_m + 1 THEN 1 ELSE 0 END) AS wrong,
         ROUND(MAX(b.tile_m - t.full_m)) AS max_extra_m
  FROM (SELECT /*+ BROADCAST(h) */ s.id, h.halo_m FROM schools s CROSS JOIN halos h) sh
  JOIN truth t ON sh.id = t.id
  LEFT JOIN best b ON sh.id = b.id AND sh.halo_m = b.halo_m
  GROUP BY sh.halo_m
),
hosp AS (
  SELECT halo_m, COUNT(*) AS hospital_rows, COUNT(DISTINCT id) AS hospitals_in_tiles
  FROM tile_hospitals GROUP BY halo_m
),
sch AS (
  SELECT halo_m, COUNT(*) AS school_rows, COUNT(*) - COUNT(DISTINCT id) AS school_duplicates
  FROM tile_schools GROUP BY halo_m
)
SELECT r.halo_m, r.schools, r.wrong, r.dropped, r.max_extra_m,
       h.hospitals_in_tiles, h.hospital_rows, s.school_rows, s.school_duplicates
FROM scored r JOIN hosp h ON r.halo_m = h.halo_m JOIN sch s ON r.halo_m = s.halo_m
ORDER BY r.halo_m
HaloWrong hospitalNo hospitalHospitals in tilesHospital rowsSchool rowsSchool duplicates
0 km121677947942,1390
5 km25527941,3903,4441,305
10 km27347941,9134,8762,737
25 km1448193,7449,4267,287
50 km008677,30118,24516,106

The query ran in 21.8 seconds. The 0 km row reproduces the tile loop: 121 and 67. Wrong answers rise from 25 to 27 between 5 and 10 km, as schools with no hospital at 5 km get one that is still not the nearest. Only the 50 km halo removes every error. Its tiles reach 73 more hospitals across the state line, and the 867 hospitals fill 7,301 rows, 8.4 copies on average. A loop reading every school in its grown tile would see the 2,139 schools as 18,245 rows, 16,106 of them repeats. The right width came from comparing against the full-extent answer, which a tile loop does not have. A sparser state, a different facility layer, or k = 3 would each need a new width.

Two bar charts over halo widths of 0, 5, 10, 25 and 50 km. Left: wrong-hospital counts 121, 25, 27, 14, 0 and no-hospital counts 67, 52, 34, 4, 0. Right: hospital rows 794, 1,390, 1,913, 3,744, 7,301 and school rows 2,139, 3,444, 4,876, 9,426, 18,245
Growing each county tile by a halo removes the seam errors only at 50 km, where the tiles hold 7,301 hospital rows for 867 hospitals and 18,245 school rows for 2,139 schools.

Features on the county line

Building footprints cross county lines. This query assigns every Overture building in a box across southeast Denver to counties two ways: by ST_Intersects on the footprint, and by the county that holds the footprint's centroid.

-- Buildings on county lines in a box across southeast Denver and Arapahoe County
WITH counties AS (
  SELECT GEOID AS county, NAME AS county_name, geometry
  FROM wherobots_open_data.us_census.tiger_county
  WHERE STATEFP = '08'
),
b AS (
  SELECT id, geometry
  FROM wherobots_open_data.overture_maps_foundation.buildings_building
  WHERE bbox.xmin BETWEEN -105.0 AND -104.88 AND bbox.ymin BETWEEN 39.64 AND 39.76
),
hits AS (
  -- Assignment by ST_Intersects: one row per building and county it touches
  SELECT b.id, c.county_name
  FROM b JOIN counties c ON ST_Intersects(c.geometry, b.geometry)
),
per AS (
  SELECT id, COUNT(*) AS n_counties, CONCAT_WS(' / ', SORT_ARRAY(COLLECT_LIST(county_name))) AS counties
  FROM hits GROUP BY id
),
cent AS (
  -- Assignment by centroid. A centroid exactly on a county line would match both counties,
  -- so production code adds a tie-break such as the lowest county GEOID.
  SELECT b.id, c.county_name
  FROM b JOIN counties c ON ST_Intersects(c.geometry, ST_Centroid(b.geometry))
),
totals AS (
  SELECT p.*,
    COUNT(*) OVER () AS buildings,
    SUM(n_counties) OVER () AS intersects_rows,
    SUM(CASE WHEN n_counties > 1 THEN 1 ELSE 0 END) OVER () AS straddling
  FROM per p
)
SELECT t.id, t.n_counties, t.counties, c.county_name AS centroid_county,
       ST_AsText(b.geometry) AS wkt,
       t.buildings, t.intersects_rows, t.straddling,
       (SELECT COUNT(*) FROM cent) AS centroid_rows
FROM totals t
JOIN b ON t.id = b.id
LEFT JOIN cent c ON t.id = c.id
WHERE t.n_counties > 1

In 15.8 seconds the query found 133,094 buildings in the box. ST_Intersects returned 133,147 building-county rows, because 53 buildings cross the line between Denver and Arapahoe County and land in both. A tile loop that keeps only buildings fully inside a tile would drop the same 53. Centroid assignment returned 133,094 rows, one per building, the same idea as the reference-point rule. A centroid exactly on the line would match both counties, so production code adds a tie-break.

Map of a box across southeast Denver with the Denver county boundary drawn in purple and 53 red markers on buildings that cross it, plus three zoomed insets of building footprints cut by the county line. Side panel: 133,094 buildings in the box, 133,147 rows from ST_Intersects, 53 buildings counted twice, 133,094 rows by centroid
53 of 133,094 Overture buildings cross the Denver and Arapahoe County line, so ST_Intersects returns 133,147 rows; centroid assignment returns one row per building.

Reading the plan

EXPLAIN FORMATTED on the nearest-hospital query showed BroadcastIndexJoin over an R-tree SpatialIndex for every point-in-polygon step, a KNNJoin for the nearest neighbor, and no CartesianProduct. The halo query planned the tile-to-school join as a partitioned RangeJoin and the nearest neighbor as BroadcastObjectSideKNNJoin, with BroadcastNestedLoopJoin only for the five-row halo list. The query states the analysis, and the plan carries the partitioning.

Spatial systems compared: what runs without a tile loop

The tax applies when the data outgrows the machine a system runs on. This table compares what each type of system provides for spatial data from its own documentation.

SystemExecutionRaster type and SQL raster functionsVector ST_ functionsImagery inference to vectorsDesigned as
PostGIS 3.6One PostgreSQL server; the Citus extension shards tables across nodesYes, PostGIS raster (161 reference entries)309Not documentedPostgreSQL extension in a transactional database
DuckDB spatialIn-process, one machineNot documented157Not documentedIn-process analytical database
GeoPandasOne Python process; dask-geopandas adds partitionsNot documentedPython API, not countedNot documentedPython library extending pandas
QGISDesktop applicationYes, raster layers and Raster CalculatorNot countedNot documentedDesktop GIS
BigQueryDistributed, serverlessThrough Earth Engine with ST_REGIONSTATSNot countedNot documentedManaged data platform
SnowflakeVirtual warehouses (compute clusters)Not documentedNot countedNot documentedCloud data platform
WherobotsDB and RasterFlowDistributed runtimes, Micro to 4X-LargeYes, 107 RS_ functions277Yes, RasterFlow (Public Preview)Spatial analytics engine for vector and raster

Sources: the PostGIS reference and PostgreSQL transactions; Citus; the DuckDB spatial function reference and DuckDB transactions; GeoPandas and dask-geopandas spatial partitioning; QGIS raster data; BigQuery overview and BigQuery raster data; Snowflake geospatial data types and virtual warehouses; the WherobotsDB reference, raster functions, runtimes, and RasterFlow. Counts are distinct ST_ function names in each official reference on 6 October 2026: PostGIS 3.6.4 core reference entries (raster, topology, and geocoder extras excluded from the vector count), the DuckDB spatial function index, and WherobotsDB reference pages. "Not documented" means the product's documentation linked here describes no such capability; "not counted" means no count was made.

Of the references counted, PostGIS lists the most vector functions (309) and WherobotsDB the next most (277), and PostGIS runs inside a transactional PostgreSQL server. DuckDB, GeoPandas, and QGIS are each bounded by one machine. BigQuery and Snowflake distribute queries over GEOGRAPHY and GEOMETRY types. WherobotsDB distributes vector and raster SQL in one engine, and RasterFlow adds mosaics, computer vision inference, and vectorized outputs on the same platform.

Partitioning tax research and what changes at scale

Edge effects, from geography to databases

Spatial statistics corrected edge effects for decades before databases met the same problem when they partitioned space to run joins. Jignesh Patel and David DeWitt's Partition based spatial-merge join (SIGMOD 1996), PBSM, partitions both inputs on a grid, copies shapes that cross an edge, joins each partition, and removes duplicates, which is the tile loop with the seams handled. Dittrich and Seeger's reference point (2000) later replaced its final duplicate-removal sort.

Spatial partitioning on clusters

Ablimit Aji and colleagues built Hadoop-GIS (PVLDB 2013) on MapReduce, handling objects that cross partition boundaries during query processing. Ahmed Eldawy and Mohamed Mokbel's SpatialHadoop (ICDE 2015) added a global spatial index. Jia Yu, Jinxuan Wu, and Mohamed Sarwat brought partitioned spatial joins to Apache Spark with GeoSpark (ACM SIGSPATIAL 2015), and the 2019 GeoInformatica paper added grid, R-tree, quadtree, and KDB-tree partitioners with local indexes. GeoSpark became Apache Sedona, and WherobotsDB comes from Sedona's original creators. On SpatialBench at scale factor 1000, the current WherobotsDB runs spatial queries and joins up to 2.5 times faster than the previous generation, with a mean speedup of 1.9 times (March 2026).

Open problems

  • Skew. KDB-tree partitions balance row counts, but cost also follows geometry complexity, so a partition holding a dense downtown or detailed coastline runs longer. Suprio Ray and colleagues traced this processing skew to the refine step in Skew-resistant parallel in-memory spatial join (SSDBM 2014).
  • Large polygons. A county, a river basin, or a country crosses many partitions and is copied to each, and every refine test walks its long boundary. Thanasis Georgiadis and Nikos Mamoulis filtered with Raster Intervals (SIGMOD 2023) before the exact test.

At national scale the same queries run on more partitions, and the seams stay inside the engine's plan, so geospatial analysis at that scale runs as one query per question.

Read more from Wherobots

Run the nearest-hospital query on your own extent with a Wherobots free trial at cloud.wherobots.com.

Frequently asked questions

What is the partitioning tax?

The partitioning tax is the work and error an analyst absorbs when a spatial dataset outgrows one machine: tiling the extent (for example by county), processing each tile, handling features on tile boundaries, reassembling the output, debugging seam artifacts, and rerunning everything when a parameter changes. Wherobots introduced the term in its Wherobots for QGIS announcement.

What is spatial partitioning?

Spatial partitioning splits a dataset into groups of rows that sit near each other in space, so each group can be stored, indexed, or processed on its own. Common schemes are a uniform grid, a quadtree, which splits crowded cells into four, and a KDB-tree, which splits at the median along alternating axes so each cell holds about the same number of rows. Distributed spatial engines such as Apache Sedona and WherobotsDB partition both sides of a join this way and handle features that cross cell edges.

What is geometric partitioning?

Geometric partitioning divides data by the coordinates of its features, so the partition key is location itself. Grids, quadtrees, k-d trees, and space-filling curves such as Hilbert and Z-order are geometric partitioning methods. Nearby features land in the same partition, which keeps spatial joins and neighbor searches local to one partition for most rows.

What are the three types of partitions?

In data systems, the three common strategies are horizontal partitioning (sharding), where each partition holds a subset of rows with the same schema; vertical partitioning, where each partition holds a subset of columns; and functional partitioning, where data is split by the part of the system that uses it. Microsoft’s data partitioning guidance describes all three. Spatial partitioning is a form of horizontal partitioning keyed on location.

What is sharding vs. partitioning?

Partitioning is the general practice of splitting a dataset into parts. Sharding is horizontal partitioning in which each part, a shard, lives in a separate data store or on a separate node, with the same schema in every shard. A spatial engine that splits a table into KDB-tree cells across workers is sharding by location for the duration of a query.

What is an edge effect in GIS?

An edge effect is a bias near the boundary of a study area or tile, where part of each feature’s neighborhood lies outside the data. Nearest-neighbor distances come out too long, point-pattern statistics such as Ripley’s K come out too low, and focal raster operations such as slope have missing neighbors. Edge corrections, buffer zones, and processing the full extent at once are the standard remedies. A tile loop creates an edge effect at every tile edge.

Is the partitioning tax the same as an edge effect?

The edge effect is the error: a bias near any boundary where part of a feature’s neighborhood lies outside the data. The partitioning tax is the full cost of a tile loop: the edge effects it creates at every tile edge, the halos and duplicate removal that repair them, and the reruns when a parameter changes. An engine that runs the full extent in one query has no internal tile edges, so only the true boundary of the data remains.

What is the modifiable areal unit problem?

The modifiable areal unit problem (MAUP) is the change in statistical results that comes from changing the zones data is aggregated into, either their size (the scale effect) or their boundaries (the zoning effect). Stan Openshaw described it in 1984, and Fotheringham and Wong showed in 1991 that it makes multivariate regression results unpredictable. Tiling by county applies one zoning to every result.

How do you avoid duplicates when tiling spatial data?

Use a deterministic rule that reports each result in exactly one tile. Assigning each feature by its centroid or a representative point, with a tie-break such as the lowest tile ID for points exactly on a tile edge, gives every feature one tile. For pairs found in more than one tile, the reference point method of Dittrich and Seeger keeps a pair only in the tile that holds a fixed corner of the intersection of the two bounding boxes. Without a rule, a final DISTINCT or GROUP BY on feature IDs removes the copies.

Why do AI coding assistants generate tile loops?

Asked to scale a spatial workflow from a county to a state on a single-machine stack, code from AI coding assistants tends to follow the familiar pattern: loop over counties, run the analysis per county, and concatenate the results. The generated loop usually has no overlap buffer and no deduplication, so it carries every seam error described here. Against an engine that handles the full extent, the same request becomes one SQL query.

Why do single-machine tools handle one large dataset but struggle with spatial joins?

A query that handles each row on its own, such as a filter, a count, or an aggregate by region, splits cleanly: each row belongs to one piece, and the pieces add up. A spatial join needs both sides of every matching pair together, and on one machine that bounds the join by that machine’s memory, or forces the data into tiles with seams. Apache Sedona’s single-node SpatialBench results show DuckDB and SedonaDB with similar low latency on spatial filters at SF 1 and SF 10, while DuckDB “encounters scaling issues in certain cases” on complex joins and lacks built-in KNN join support. WherobotsDB distributes the join across a cluster, so the same SQL runs on the full extent.

Does Wherobots need a tile loop for a large extent?

No. WherobotsDB runs one query across the full extent and handles spatial partitioning, partition-level indexes, copies of features that cross partition edges, and duplicate removal inside the engine. When the extent grows, pick a larger runtime; Wherobots runtimes range from General Purpose Micro to 4X-Large, and the query text stays the same.