Skip to content

Commit bd20e35

Browse files
committed
chg: use Duckdb Python API for GBIF backbone
1 parent 6480d4d commit bd20e35

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,8 +1,8 @@
11
import re
22

3+
from duckdb import ConstantExpression, DuckDBPyConnection, FunctionExpression
34
from lxml import etree
45
from lxml.etree import _Element
5-
from pyarrow import csv
66

77
PARSER: etree.XMLParser = etree.XMLParser(resolve_entities=False)
88

@@ -32,13 +32,61 @@ def get_anytext(bag: str | _Element | list[str]) -> str:
3232
raise TypeError("xpath result was not a list of strings")
3333

3434

35-
TSV_PARSE_OPTIONS = csv.ParseOptions(
36-
delimiter="\t", quote_char=False, double_quote=False
37-
)
35+
class DuckDBAtomicTransaction:
36+
def __init__(self, conn: DuckDBPyConnection):
37+
self.conn = conn
3838

39+
def __enter__(self):
40+
self.conn.begin()
41+
return self.conn
3942

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

0 commit comments

Comments
 (0)