Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
50 changes: 44 additions & 6 deletions task/pgload/load.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,13 +20,15 @@

DATABASE_NAME = os.getenv("DATABASE_NAME")
logger = logging.getLogger(__name__)
logger.setLevel(logging.INFO)
logger.propagate = False
logging.basicConfig(level=logging.INFO)
streamHandler = logging.StreamHandler(sys.stdout)
formatter = logging.Formatter("%(asctime)s - %(name)s - %(levelname)s - %(message)s")
streamHandler.setFormatter(formatter)
logger.addHandler(streamHandler)

complex_geom_datasets = ['flood-risk-zone']
complex_geom_datasets = ["flood-risk-zone"]
export_tables = {
DATABASE_NAME: ["entity", "old_entity"],
"digital-land": [
Expand Down Expand Up @@ -126,8 +128,10 @@ def do_replace_table(table, source, csv_filename, postgress_conn, sqlite_conn):

make_valid_with_handle_geometry_collection(postgress_conn, source)

null_invalid_geometries(postgress_conn, source)

if source in complex_geom_datasets:
update_entity_subdivided(postgress_conn,source)
update_entity_subdivided(postgress_conn, source)


def do_replace(source, sqlite_conn, tables_to_export=None):
Expand Down Expand Up @@ -225,11 +229,43 @@ def make_valid_multipolygon(connection, source):

logger.info(f"Updated {rowcount} rows with valid multi polygons")

def update_entity_subdivided(connection, source):

def null_invalid_geometries(connection, source):
"""
Null the geometry and geojson of any entity whose geometry is still invalid
after the PostGIS repair step. The entity row is kept — only the unusable
shape is dropped so it stops breaking tiles and spatial queries.

This is non-destructive and re-derived on every load: the invalid data
remains in the source resources, and the pipeline already raises an
'invalid geometry - not fixable' issue (which surfaces as an LPA task), so
no issue is written here.
"""
null_invalid = """
UPDATE entity
SET geometry = NULL,
geojson = NULL
WHERE dataset = %s
AND geometry IS NOT NULL
AND NOT ST_IsValid(geometry);
""".strip()

with connection.cursor() as cursor:
cursor.execute(null_invalid, (source,))
rowcount = cursor.rowcount
connection.commit()

logger.info(
f"null_invalid_geometries: nulled {rowcount} invalid "
f"geometries for dataset '{source}'"
)


def update_entity_subdivided(connection, source):
delete_sql = """
DELETE FROM entity_subdivided WHERE dataset = %s ;
"""

update_entity_subdivided = """
INSERT INTO entity_subdivided (entity, dataset, geometry_subdivided)
SELECT e.entity, e.dataset, ST_Multi(g.geom)
Expand All @@ -244,7 +280,6 @@ def update_entity_subdivided(connection, source):
""".strip()

with connection.cursor() as cursor:

cursor.execute(delete_sql, (source,))
deleted_count = cursor.rowcount

Expand All @@ -254,7 +289,10 @@ def update_entity_subdivided(connection, source):
connection.commit()

logger.info(f"Updated entity_sub_divided table - {deleted_count} rows deleted")
logger.info(f"Updated entity_sub_divided table - {rowcount} rows with subdivided geometries")
logger.info(
f"Updated entity_sub_divided table - {rowcount} rows with subdivided geometries"
)


if __name__ == "__main__":
do_replace_cli()
68 changes: 66 additions & 2 deletions tests/integration/pg_load/test_load.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
from task.pgload.load import ( # noqa: E402
make_valid_multipolygon,
make_valid_with_handle_geometry_collection,
null_invalid_geometries,
SQL,
call_sql_queries,
export_tables,
Expand Down Expand Up @@ -112,6 +113,65 @@ def test_make_valid_with_handle_geometry_collection(postgresql_conn, sources):
cursor.close()


def test_null_invalid_geometries(postgresql_conn, create_db):
source = "article-4-direction"
invalid_entity = 9999991
valid_entity = 9999992

cursor = postgresql_conn.cursor()
# unfixable (self-intersecting bowtie) with a geojson value alongside it
cursor.execute(
"""
INSERT INTO entity (entity, name, dataset, reference, typology, geometry, geojson)
VALUES (
%s, 'Unfixable invalid geometry', %s, 'unfixable-invalid-geometry',
'geography',
ST_GeomFromText('POLYGON((0 0,1 1,1 0,0 1,0 0))'),
'{"type":"Polygon"}'
);
""",
(invalid_entity, source),
)
# a valid geometry that must be left untouched
cursor.execute(
"""
INSERT INTO entity (entity, name, dataset, reference, typology, geometry, geojson)
VALUES (
%s, 'Valid geometry', %s, 'valid-geometry', 'geography',
ST_GeomFromText('POLYGON((0 0,0 1,1 1,1 0,0 0))'),
'{"type":"Polygon"}'
);
""",
(valid_entity, source),
)
postgresql_conn.commit()
cursor.close()

null_invalid_geometries(postgresql_conn, source)

cursor = postgresql_conn.cursor()
# invalid entity: row kept, geometry + geojson nulled, other fields intact
cursor.execute(
"SELECT geometry, geojson, name, reference FROM entity WHERE entity = %s",
(invalid_entity,),
)
geometry, geojson, name, reference = cursor.fetchone()
assert geometry is None
assert geojson is None
assert name == "Unfixable invalid geometry"
assert reference == "unfixable-invalid-geometry"

# valid entity untouched
cursor.execute(
"SELECT geometry IS NOT NULL, geojson IS NOT NULL FROM entity WHERE entity = %s",
(valid_entity,),
)
geometry_present, geojson_present = cursor.fetchone()
assert geometry_present
assert geojson_present
cursor.close()


def test_unretired_entities(postgresql_conn):
source = "certificate-of-immunity"
table = "old_entity"
Expand Down Expand Up @@ -171,12 +231,16 @@ def test_remove_invalid_datasets_deletes_from_old_entity(postgresql_conn, create
cursor.close()

with patch("task.pgload.load.get_pg_connection", return_value=postgresql_conn):
remove_invalid_datasets(["valid-dataset", "article-4-direction", "ancient-woodland"])
remove_invalid_datasets(
["valid-dataset", "article-4-direction", "ancient-woodland"]
)

cursor = postgresql_conn.cursor()

cursor.execute("SELECT COUNT(*) FROM old_entity WHERE dataset = 'invalid-dataset'")
assert cursor.fetchone()[0] == 0, "invalid-dataset rows should be removed from old_entity"
assert (
cursor.fetchone()[0] == 0
), "invalid-dataset rows should be removed from old_entity"

cursor.execute("SELECT COUNT(*) FROM old_entity WHERE old_entity = 9990001")
assert cursor.fetchone()[0] == 1, "valid-dataset row should remain in old_entity"
Expand Down