Skip to content

Commit ae923a4

Browse files
committed
chg: use Duckdb Python API for GBIF backbone
1 parent d9febee commit ae923a4

4 files changed

Lines changed: 99 additions & 43 deletions

File tree

src/datasync/gbif_backbone.py

Lines changed: 43 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -1,20 +1,17 @@
1-
from importlib.resources import files
2-
3-
import duckdb
41
import fsspec
52
import typer
3+
from duckdb import (
4+
DuckDBPyConnection,
5+
StarExpression,
6+
connect,
7+
)
68

7-
from .libs.helpers import pa_read_tsv
9+
from .libs.helpers import DuckDBAtomicTransaction, create_fts_index, function_expression
810
from .settings import (
911
env,
1012
log,
1113
)
1214

13-
resources = files(__package__).joinpath("gbif_backbone")
14-
taxon_sql = resources.joinpath("taxon.sql").read_text()
15-
vernacularname_sql = resources.joinpath("vernacularname.sql").read_text()
16-
17-
1815
log.debug("Importing GBIF Backbone settings")
1916

2017

@@ -29,23 +26,51 @@
2926
app = typer.Typer(help="export GBIF Backbone data to DuckDB database")
3027

3128

32-
def import_taxon(conn, archive):
29+
def import_taxon(conn: DuckDBPyConnection, archive):
3330
log.debug("Importing Taxon.tsv")
34-
pa_taxon = pa_read_tsv(archive, "Taxon.tsv") # noqa: F841
35-
conn.execute(taxon_sql)
31+
conn.sql("DROP TABLE IF EXISTS taxon")
32+
conn.from_csv_auto(archive.open("Taxon.tsv")).to_table("taxon")
33+
conn.sql("CREATE INDEX taxon_taxonid ON taxon (taxonid)")
34+
conn.sql(
35+
create_fts_index(
36+
"taxon",
37+
"taxonID",
38+
"canonicalName",
39+
overwrite=True,
40+
)
41+
)
3642

3743

38-
def import_vernacular_name(conn, archive):
44+
def import_vernacular_name(conn: DuckDBPyConnection, archive):
3945
log.debug("Importing VernacularName.tsv")
40-
pa_vernacular_names = pa_read_tsv(archive, "VernacularName.tsv") # noqa: F841
41-
conn.execute(vernacularname_sql)
46+
conn.sql("DROP TABLE IF EXISTS vernacular_name")
47+
generate_vernacular_id = function_expression(
48+
"concat_ws",
49+
"|",
50+
"taxonID",
51+
"language",
52+
"vernacularName",
53+
)
54+
(
55+
conn.from_csv_auto(archive.open("VernacularName.tsv"))
56+
.select(StarExpression(), generate_vernacular_id.alias("vernacularID"))
57+
.to_table("vernacular_name")
58+
)
59+
conn.sql("CREATE INDEX vernacular_name_taxonid ON vernacular_name (taxonid)")
60+
conn.sql("CREATE INDEX vernacular_name_language ON vernacular_name (language)")
61+
conn.sql(
62+
create_fts_index(
63+
"vernacular_name", "vernacularID", "vernacularName", overwrite=True
64+
)
65+
)
4266

4367

4468
@app.command()
4569
def import_all():
4670
"""Import GBIF Backbone data into a DuckDB database."""
4771
archive = fsspec.filesystem("zip", fo=GBIF_BACKBONE_URL.geturl(), mode="r")
48-
with duckdb.connect(GBIF_BACKBONE_DUCKDB_NAME) as conn:
49-
import_taxon(conn, archive)
50-
import_vernacular_name(conn, archive)
72+
with connect(GBIF_BACKBONE_DUCKDB_NAME) as conn:
73+
with DuckDBAtomicTransaction(conn):
74+
import_taxon(conn, archive)
75+
import_vernacular_name(conn, archive)
5176
log.info("GBIF Backbone data imported successfully")

src/datasync/gbif_backbone/taxon.sql

Lines changed: 0 additions & 6 deletions
This file was deleted.

src/datasync/gbif_backbone/vernacularname.sql

Lines changed: 0 additions & 11 deletions
This file was deleted.

src/datasync/libs/helpers.py

Lines changed: 56 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
import re
22

3+
from duckdb import ConstantExpression, DuckDBPyConnection, FunctionExpression
34
from lxml import etree
4-
from pyarrow import csv
55

66
PARSER = etree.XMLParser(resolve_entities=False)
77

@@ -26,13 +26,61 @@ def get_anytext(bag: str) -> str:
2626
)
2727

2828

29-
TSV_PARSE_OPTIONS = csv.ParseOptions(
30-
delimiter="\t", quote_char=False, double_quote=False
31-
)
29+
class DuckDBAtomicTransaction:
30+
def __init__(self, conn: DuckDBPyConnection):
31+
self.conn = conn
3232

33+
def __enter__(self):
34+
self.conn.begin()
35+
return self.conn
3336

34-
def pa_read_tsv(archive, filename):
35-
return csv.read_csv(
36-
archive.open(filename),
37-
parse_options=TSV_PARSE_OPTIONS,
37+
def __exit__(self, exc_type, exc_value, traceback):
38+
if exc_type is None:
39+
self.conn.commit()
40+
else:
41+
self.conn.rollback()
42+
43+
44+
def parameter_expression(value):
45+
return (
46+
str(FunctionExpression("", ConstantExpression("").alias(value)))
47+
.removeprefix("(")
48+
.removesuffix(" := '')")
49+
)
50+
51+
52+
def identity_expression(value):
53+
"""SQLIdentifier is not exposed by DuckDB https://github.com/duckdb/duckdb/discussions/20425"""
54+
return (
55+
str(FunctionExpression("", ConstantExpression(value).alias("")))
56+
.removeprefix("(")
57+
.removesuffix(")")
58+
)
59+
60+
61+
def named_argument_expression(name, value):
62+
return parameter_expression(name) + " = " + identity_expression(value)
63+
64+
65+
def function_expression(name, *args):
66+
expressions = [ConstantExpression(arg) for arg in args]
67+
function = FunctionExpression(name, *expressions)
68+
return function
69+
70+
71+
def function_expression_with_named_arguments(name, *args, **kwargs):
72+
function = function_expression(name, *args)
73+
named_arguments = []
74+
for key, value in kwargs.items():
75+
named_argument = named_argument_expression(key, value)
76+
named_arguments.append(named_argument)
77+
return str(function)[:-1] + ", " + ", ".join(named_arguments) + ")"
78+
79+
80+
def create_fts_index(*args, **kwargs):
81+
expression = function_expression_with_named_arguments(
82+
"create_fts_index",
83+
*args,
84+
**kwargs,
3885
)
86+
return f"PRAGMA {expression}"

0 commit comments

Comments
 (0)