Quick start¶
This example searches arXiv, asks an LLM which papers are relevant, extracts
measurements from their full text, and upserts them into a CSV. It needs the
async, llm, and pdf extras, and reads the API key from the LLM_API_KEY
environment variable, so the key never appears in source code. To load it from
a .env file instead, see Configuration.
import asyncio
import os
from sci_etl_core import (
AsyncArxivExtractor,
AsyncCsvUpsertExporter,
AsyncETLPipeline,
AsyncFileStateManager,
AsyncLLMEntityExtractor,
AsyncLLMRelevanceFilter,
AsyncOpenAICompatibleClient,
PipelineAborted,
)
from sci_etl_core.http_async import build_async_client
from sci_etl_core.parsers import LatexTarballParser, PdfPlumberParser
from sci_etl_core.processors import DefaultKeyNormalizer
RELEVANCE_PROMPT = (
"Decide whether the paper reports measurements of galaxies. "
'Reply with JSON: {"relevant": true} or {"relevant": false}.'
)
EXTRACTION_PROMPT = (
"Extract every measured object from the paper. Reply with JSON: "
'{"items": [{"name": "...", "value_a": 0.0, "value_b": 0.0}]}.'
)
async def main() -> None:
client = build_async_client()
llm = AsyncOpenAICompatibleClient(
api_key=os.environ["LLM_API_KEY"],
base_url="https://api.openai.com/v1",
model="gpt-4o-mini",
)
pipeline = AsyncETLPipeline(
extractor=AsyncArxivExtractor(
client=client,
pdf_parser=PdfPlumberParser(),
latex_parser=LatexTarballParser(),
),
relevance_filter=AsyncLLMRelevanceFilter(
llm_client=llm, system_prompt=RELEVANCE_PROMPT
),
entity_extractor=AsyncLLMEntityExtractor(
llm_client=llm, system_prompt=EXTRACTION_PROMPT
),
exporter=AsyncCsvUpsertExporter(
key_column="name",
value_columns=["value_a", "value_b"],
normalizer=DefaultKeyNormalizer(),
),
state_manager=AsyncFileStateManager("state/processed.txt", "state/metadata.json"),
destination="results.csv",
max_concurrency=4,
logger=print,
closeables=[client, llm],
)
async with pipeline:
try:
processed = await pipeline.run(query="all:galaxy", total_limit=50)
except PipelineAborted as exc:
print(f"Stopped early after {exc.partial_count} records: {exc}")
return
print(f"Processed {processed} relevant records")
asyncio.run(main())
What a run does¶
- Fetches a listing page (
page_sizerecords, default 100) starting at the offset saved by the state manager (or atstart_index=, if given), and skips records already processed. A page that holds only processed records is passed over, not treated as the end of the data. - For each remaining record, with at most
max_concurrencyin flight: relevance filter → full-text fetch → optional memory ingest → entity extraction → export → mark processed. Irrelevant records are marked processed without fetching full text. Records whoserecord_idis missing or blank can't be tracked, so they are skipped and logged. - Saves the listing offset and repeats until
total_limitrelevant records are processed or the listing is empty, waitingsleep_betweenseconds (default 0) before each further page. The offset only moves past pages whose records were all settled; see State, resuming, and errors.
total_limit counts relevant records only, is never exceeded, and defaults
to page_size. max_concurrency and page_size must be at least 1 and
total_limit must not be negative; other values raise ValueError before any
request is made. The deprecated max_records= argument sets both values
and will be removed in 0.5.
AsyncArxivExtractor also waits sleep_before_search seconds (default 3)
before every listing request, to respect arXiv's rate limits.
On exit, async with pipeline awaits aclose() on every entry in
closeables that has one — the HTTP client, the LLM client, and any
AsyncSqliteStateManager, AsyncSqliteEmbeddingStore, or
AsyncSqliteFts5Store you use.
Prompts must ask for JSON¶
AsyncOpenAICompatibleClient requests JSON mode
(response_format={"type": "json_object"}), and OpenAI's API rejects JSON-mode
requests whose messages never mention "JSON". The response shapes the library
reads:
AsyncLLMRelevanceFilterreads therelevantkey. It accepts a boolean,0/1, or the strings"true","false","yes","no","1", and"0"in any case. Anything else, including a missing key, is treated as an error and returnsdefault_on_error.AsyncLLMEntityExtractorreads the list underresult_key(default"items"), or the only value if the response has exactly one key. The list must hold objects;nullmeans no entities and a lone object counts as one. Any other value there raisesLLMError, so the record is retried. A response with several keys but noresult_keyyields no entities, and the record is marked processed, so name the key in the prompt.AsyncCsvUpsertExportertakes each item'skey_columnvalue as the row key, so the extraction prompt must ask for that field.
Next steps¶
- Run the same pipeline from synchronous code with
ETLPipeline. - Load the settings from YAML and
.envwith typed configuration. - Clean and plot the CSV with the post-processing steps.