A destination receives each done job as a file, and a run log records what happened to each payload. Together they let a large run write its results to disk or a bucket, and leave a record to recover from.
Write each job to a destination
destination names a folder. The run writes each done job there as <job_id>.json before it yields the job, so a yielded job’s file always exists:
import os
import oxyscraper as oxy
USERNAME = os.environ["OXY_WSA_USERNAME"]
PASSWORD = os.environ["OXY_WSA_PASSWORD"]
payloads = [
oxy.Universal(url=f"https://sandbox.oxylabs.io/products/{number}")
for number in range(1, 4)
]
with oxy.Session(username=USERNAME, password=PASSWORD) as session:
run = session.execute(payloads, destination="results")
for _ in run:
pass
print(run.progress)
sorted(os.listdir("results"))
3/3 done, 3 written, 1.00s
['7500000000000000001.json',
'7500000000000000002.json',
'7500000000000000003.json']
Drain such a run with a bare loop, as above. all() returns every job with its content, so it holds the whole run in memory.
Each file holds the API’s body unchanged, as one line of JSON. Only done jobs get a file, so every file in a destination is a billed result. A run never deletes a file, and replaces a file of the same name, so you choose what to clear between runs.
destination takes three kinds of value:
- A local path, which oxy creates.
- A URL such as
s3://bucket/path, gs://bucket/path or az://container/path, which reads the same AWS_*, GOOGLE_* and AZURE_* variables as Polars.
- An obstore store, for explicit credentials, a credential provider, an S3 region or a retry config.
The run lists the destination before its first submission, so a missing bucket or a wrong credential raises before anything bills. A failed write stops submission, and the run yields the remaining jobs unwritten before it raises.
Read a destination with Polars
pl.scan_ndjson reads the folder. Pass a schema, because the values in job.context mix lists and scalars, and Polars cannot infer one type for them:
import polars as pl
schema = {
"results": pl.List(
pl.Struct({"url": pl.String, "status_code": pl.Int64, "content": pl.String})
),
"job": pl.Struct({"id": pl.String, "source": pl.String, "status": pl.String}),
}
(
pl.scan_ndjson("results/*.json", schema=schema)
.explode("results", empty_as_null=False)
.select(
pl.col("job").struct.field("id"),
pl.col("results").struct.field("url", "status_code"),
)
.collect()
)
shape: (3, 3)| id | url | status_code |
|---|
| str | str | i64 |
| "7500000000000000001" | "https://sandbox.oxylabs.io/pro… | 200 |
| "7500000000000000002" | "https://sandbox.oxylabs.io/pro… | 200 |
| "7500000000000000003" | "https://sandbox.oxylabs.io/pro… | 200 |
A schema also limits the scan to the fields it names. content is a string for raw and markdown results, and an object for parsed ones, so give it the type of the output type you scan.
Write a run log
run_log names a folder, and takes the same values as destination. Each run writes one file there, named by the run’s start in UTC, such as 20260929T170412.345Z.jsonl:
with oxy.Session(username=USERNAME, password=PASSWORD) as session:
run = session.execute(payloads, destination="results", run_log="logs")
try:
for _ in run:
pass
except oxy.IncompleteRunError as error:
print(error)
Stopped checking job 7500000000000000001, because the API returned 404 Not Found: universal https://sandbox.oxylabs.io/products/1
Job 7500000000000000002 faulted: universal https://sandbox.oxylabs.io/products/2
The file holds one line per payload, in input order, and each line has the same five keys:
import json
from pathlib import Path
log = max(Path("logs").glob("*.jsonl"))
lines = [json.loads(line) for line in log.read_text().splitlines()]
lines
[{'state': 'unfetched',
'id': '7500000000000000001',
'payload': {'source': 'universal',
'url': 'https://sandbox.oxylabs.io/products/1'},
'error': {'status_code': 404,
'message': 'Not Found',
'trace_id': '00000001-000000000000000000000000'},
'upload': None},
{'state': 'faulted',
'id': '7500000000000000002',
'payload': {'source': 'universal',
'url': 'https://sandbox.oxylabs.io/products/2'},
'error': None,
'upload': None},
{'state': 'done',
'id': '7500000000000000003',
'payload': {'source': 'universal',
'url': 'https://sandbox.oxylabs.io/products/3'},
'error': None,
'upload': None}]
state is done, faulted, rejected, unsubmitted or unfetched.
id is the job’s ID, or null when no job exists.
payload is the body that oxy sent.
error holds the error behind a rejected, unsubmitted or unfetched payload.
upload holds the Cloud Storage entry’s code and message.
The run writes an empty file before its first submission, so a folder it cannot write to raises before anything bills. It replaces the file when the run ends, raises or stops, even after Ctrl+C. An empty run log therefore means a process killed outright or a second Ctrl+C, and jobs may have billed without a record.
Recover a run
Nothing resumes a run, and a second call with the same payloads submits each one again. The run log lists what to redo instead. Resubmit the faulted payloads, and fetch the unfetched jobs by ID, so you pay for neither twice:
faulted = [
oxy.Payload.model_validate(line["payload"])
for line in lines
if line["state"] == "faulted"
]
unfetched = [line["id"] for line in lines if line["state"] == "unfetched"]
with oxy.Session(username=USERNAME, password=PASSWORD) as session:
for job in session.execute(faulted, destination="results"):
print(job.status, job.input)
for job_id in unfetched:
job = session.get(job_id)
print(job.status, job.input)
done https://sandbox.oxylabs.io/products/2
done https://sandbox.oxylabs.io/products/1
get writes nothing to a destination, and oxy get -d does.
The run log replaces the credentials in a storage_url with redacted:redacted. So a tos or s3_compatible payload from a run log needs its secret again before a resubmission.
oxyscraper is not affiliated with or endorsed by Oxylabs. Oxylabs and Oxy are trademarks of Oxylabs.