Skip to content

Commit 4e4aa12

Browse files
taheeraahmedfrafra
andcommitted
chg(coat): coat pipeline
Co-authored-by: Francesco Frassinelli <francesco.frassinelli@nina.no>
1 parent 45a4a9a commit 4e4aa12

2 files changed

Lines changed: 196 additions & 0 deletions

File tree

src/datasync/coat.py

Lines changed: 194 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,194 @@
1+
import json
2+
from datetime import datetime
3+
from urllib.parse import urlparse
4+
5+
import dlt
6+
import typer
7+
from dlt.sources.helpers.rest_client import RESTClient
8+
from dlt.sources.helpers.rest_client.paginators import OffsetPaginator
9+
10+
from .settings import log
11+
12+
app = typer.Typer(help="Export COAT CKAN packages to Parquet")
13+
14+
PAGE_SIZE = 100
15+
BASE_URL = "https://data.coat.no"
16+
ENDPOINT_URL = BASE_URL + "/api/3/action"
17+
18+
MIXED = ("author_email", "maintainer", "maintainer_email")
19+
SPLIT = (
20+
"associated_parties",
21+
"datasets",
22+
"funding",
23+
"location",
24+
"persons",
25+
"protocol",
26+
"scientific_name",
27+
)
28+
DATE_FIELDS = ("embargo", "temporal_start", "temporal_end")
29+
30+
31+
def normalize_record(record, domain: str = BASE_URL):
32+
"""Normalize a package record in-place."""
33+
org = record.get("organization") or {}
34+
35+
# Handle mixed fields (can be string or list)
36+
for field in MIXED:
37+
value = record.get(field)
38+
if isinstance(value, list):
39+
record[field] = json.dumps(value)
40+
elif isinstance(value, str) and value.startswith("["):
41+
try:
42+
parsed = json.loads(value)
43+
if isinstance(parsed, list) and len(parsed) == 1:
44+
record[field] = parsed[0]
45+
except json.JSONDecodeError:
46+
log.error(f"Failed to parse JSON for field {field}: {value}")
47+
raise ValueError(
48+
f"Failed to parse JSON for field {field}: {value}"
49+
) from None
50+
51+
# Split comma-separated fields into arrays
52+
for field in SPLIT:
53+
value = record.get(field)
54+
record[field] = [v.strip() for v in value.split(",")] if value else []
55+
56+
# Parse date fields
57+
for field in DATE_FIELDS:
58+
value = record.get(field)
59+
if value:
60+
try:
61+
record[field] = datetime.strptime(value, "%Y-%m-%d").date()
62+
except (ValueError, TypeError):
63+
pass
64+
65+
# Extract organization info
66+
record["organization_name"] = org.get("name")
67+
record["organization_title"] = org.get("title")
68+
69+
# Flatten extras to JSON object
70+
record["extras"] = {e["key"]: e["value"] for e in record.get("extras", [])}
71+
72+
# Extract tag names
73+
record["tags"] = [t["name"] for t in record.get("tags", [])]
74+
75+
# Build URL
76+
name = record.get("name")
77+
if not name:
78+
log.warning(
79+
f"Record {record.get('id')!r} is missing 'name', URL will be incomplete"
80+
)
81+
record["url"] = f"{domain}/dataset/{name or ''}"
82+
83+
return record
84+
85+
86+
@dlt.resource(
87+
name="packages",
88+
primary_key="id",
89+
write_disposition="replace",
90+
)
91+
def packages(
92+
api_key: str = "",
93+
base_url: str = ENDPOINT_URL,
94+
):
95+
"""Fetch packages from CKAN API and normalize them."""
96+
log.info(f"Starting package extraction from {base_url}")
97+
98+
parsed = urlparse(base_url)
99+
if not parsed.scheme or not parsed.netloc:
100+
log.error("Invalid base_url: %s", base_url)
101+
raise ValueError(f"Invalid base_url: {base_url!r}")
102+
domain = f"{parsed.scheme}://{parsed.netloc}"
103+
104+
client = RESTClient(
105+
base_url=base_url,
106+
headers={"Authorization": api_key, "Accept": "application/json"},
107+
)
108+
109+
paginator = OffsetPaginator(
110+
limit=PAGE_SIZE,
111+
offset_param="start",
112+
limit_param="rows",
113+
total_path="result.count",
114+
)
115+
116+
total_records = 0
117+
for page in client.paginate(
118+
"/ckan_package_search",
119+
params={"include_private": "true"},
120+
paginator=paginator,
121+
):
122+
for record in page:
123+
yield normalize_record(record, domain=domain)
124+
total_records += 1
125+
126+
log.info(f"Extraction complete. Total records: {total_records}")
127+
128+
129+
@dlt.transformer(
130+
data_from=packages,
131+
name="resources",
132+
primary_key="id",
133+
write_disposition="replace",
134+
)
135+
def resources(pkg):
136+
"""Extract resources from packages."""
137+
for res in pkg.get("resources", []):
138+
yield {
139+
**res,
140+
"package_name": pkg.get("name"),
141+
"package_id": pkg.get("id"),
142+
}
143+
144+
145+
@dlt.source(name="coat", max_table_nesting=0)
146+
def coat_source(api_key: str = "", base_url: str = ENDPOINT_URL):
147+
"""Define the COAT data source."""
148+
return packages(api_key=api_key, base_url=base_url), resources
149+
150+
151+
@app.command()
152+
def get_packages_and_resources(
153+
bucket_url: str = typer.Option(
154+
default="data",
155+
envvar="COAT_BUCKET_URL",
156+
help="Destination bucket URL (local path or s3://...)",
157+
),
158+
api_key: str = typer.Option(
159+
...,
160+
envvar="COAT_API_KEY",
161+
help="CKAN API key",
162+
),
163+
base_url: str = typer.Option(
164+
default="https://data.coat.no/api/3/action",
165+
envvar="COAT_BASE_URL",
166+
help="COAT CKAN API base URL",
167+
),
168+
dataset_name: str = typer.Option(
169+
default="coat",
170+
envvar="COAT_DATASET_NAME",
171+
help="Dataset name for the extracted data",
172+
),
173+
):
174+
"""Run the COAT extraction pipeline."""
175+
pipeline = dlt.pipeline(
176+
pipeline_name=dataset_name,
177+
destination=dlt.destinations.filesystem(
178+
bucket_url=bucket_url, layout="{table_name}.{ext}"
179+
),
180+
dataset_name=dataset_name,
181+
)
182+
run = pipeline.run(
183+
coat_source(api_key=api_key, base_url=base_url), loader_file_format="parquet"
184+
)
185+
log.debug(
186+
f"COAT pipeline output written to:\n"
187+
f"- {bucket_url}/coat/resources.parquet\n"
188+
f"- {bucket_url}/coat/packages.parquet"
189+
)
190+
log.info(f"Pipeline run completed. Load info: {run}")
191+
192+
193+
if __name__ == "__main__":
194+
app()

src/datasync/main.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
import typer
66

77
from . import (
8+
coat,
89
dms,
910
gbif_backbone,
1011
grass,
@@ -19,6 +20,7 @@
1920
app = typer.Typer(
2021
help="Provide subcommands for synchronizing different resources, see subcommands"
2122
)
23+
app.add_typer(coat.app, name="coat")
2224
app.add_typer(nva.app, name="nva")
2325
app.add_typer(ubw.app, name="ubw")
2426
app.add_typer(dms.app, name="dms")

0 commit comments

Comments
 (0)