Eve Files and ETL
Csv, Json, Dat), reading and writing in Unicode, rejects, and the pattern of a restartable transfer.spec/semantics/data-types.md and spec/library/data-language.md. The database side is on the page Databases.The cycle
files --read--> records --transform--> stage (SQLite) --transfer--> remote tables
^ |
+-----------write <-- records <------ report (query) <-----------------+
Each step is a job, so each step is a transaction and a place to restart.
File types
A value of a file type is made by opening a file or a stream of bytes (for example the body of an HTTP answer). Opening reads only the first block. The parser hands out the units of the block one by one, so the memory used does not depend on the size of the file: a 20 GB file needs the memory of one batch.
| Type | The file | A for loop gives | Version |
|---|---|---|---|
Csv | separated fields, usually with a header row | a row | 0.4 |
Json | JSON text or JSON Lines | a node: an object, an array, a scalar | 0.4 |
Dat | fixed-width fields | a row | 0.4 |
Xml | XML text | an element | 0.5 |
Html, Htmlt | HTML, and HTML templates | a node | 0.6, with the templates |
All of them are classes derived from Document, which is an Iterable, so a for loop reads any of them. A unit that is not well formed raises a FormatError with the line and the column when the loop reaches it, and the units before it were already given to the loop.
Untyped reading
The first way reads rows and leaves every value as text. A row is read by name or by position:
from "lib" use (csv);
process main is
new orders: Csv := csv.open("orders.csv"); ** nothing is read yet but the first block
new sum := 0.0;
for row in orders do ** one row at a time
let sum += row["total"].real(); ** row[2] reads by position
done;
print "total = {sum % f10.2}";
return;
Typed reading
The second way uses the record class as the layout of the file, the same class that is the layout of a table. The values are converted to the field types, and row.name becomes o.name:
class Order = {id: Integer, customer_id: Integer, amount: Decimal, created: Instant} <: Record;
new rows := csv.read!(:Order)("in/orders.csv", encoding: "windows-1252",
columns: {"Order No": "id", "Total": "amount"});
new docs := json.read!(:Order)("in/orders.jsonl"); ** JSON Lines or an array of objects
new lines := dat.read!(:Order)("in/orders.dat", widths: (8, 10, 12, 20));
for o in rows do
print "{o.id}: {o.amount}";
done;
- A CSV header name or a JSON key maps to the field of the same name; the option
columns:renames, ascolumn ... asdoes for a table; - A
Datfile maps by position: the widths of the fields, in the order of the fields of the class; - A missing value of a field
T?isnull; a missing value of a fieldTis an error of that row; - The header is checked at the first row: in debug mode every difference is listed, in production the first one stops the read;
- A read of a file carries
!, like every read of the outside world;writehas none.
The source can be a path or a stream of bytes, for example the body of an answer of the HTTP client. A path that ends in .gz is decompressed while it is read.
Unicode streams
Every file is read as bytes, decoded to code points, then parsed into units. The rules are the same for every format, and nothing is guessed from the content:
| Point | Rule |
|---|---|
| Encoding | a BOM decides (UTF-8, UTF-16, UTF-32); else the option encoding:; else UTF-8. ISO-8859-1 and Windows-1252 are also known |
| Invalid bytes | invalid: "error" (default) raises FormatError with line, column and byte offset; "replace" writes U+FFFD and counts it; "reject" sends the unit to the rejects |
| Normalization | normalize: "nfc" composes the text, so that two spellings of the same name are one key. Recommended for key fields; off by default |
| Line endings | LF, CRLF and CR are read; LF is written unless eol: "crlf" |
| Fixed width | widths count code points; width_unit: "byte" for old single-byte files |
Transform and rejects
A transform is a pure function from one record to another. Pure functions can run in the parallel aspects of a transfer without any lock:
function to_target(o: Order) => (@t: Target) is
let t := {id: o.id, total: o.amount * 1.19d, day: o.created.date()} :Target;
return;
A line that can not be read or converted, or a record whose validate fails, does not have to stop the whole job. Give the read a rejects output. The rejects are a table of the source, line, column, text and reason; max_rejects turns the one after the limit into an error of the job:
new bad := Table(:Reject)();
new rows := csv.read!(:Order)(path, rejects: @bad, max_rejects: 100);
...
csv.write("out/orders_rejects.csv", bad) if bad.count() > 0;
Every count (read, converted, rejected, loaded) goes to the run log, as the lineage of the run.
Reports: the other direction
A report reads a stream from a query and writes a file with the same record class and the same options. The writer takes a stream, so a report of millions of rows needs the memory of one batch. A report job sets $read_only = True for one consistent snapshot without locks:
new rows := sales.query!(:Order)("select id, customer_id, amount, created from orders where created >= ?", since);
csv.write("out/orders.csv", rows, encoding: "utf-8", bom: False, eol: "lf");
json.write("out/orders.jsonl", rows);
dat.write("out/orders.dat", rows, widths: (8, 10, 12, 20));
bom: True is for spreadsheets that need it.
Canonical form: a round trip without loss
The same class and the same options give a lossless round trip. read(write(rows)) gives the same records, and write(read(file)) gives the same file when the file is in canonical form, which every writer produces:
| Value | Written as |
|---|---|
Integer | decimal digits, - when negative |
Decimal | the digits with the scale of the value, 12.50, never an exponent |
Real | the shortest text that reads back to the same bits |
Logic | true, false (logic: ("Y", "N") changes it) |
Instant | ISO 8601 in UTC with Z; a DateTime keeps its offset |
null | CSV: an empty field without quotes ("" is the empty string, so the two stay different); JSON: null; Dat: spaces |
| CSV quotes | only when the value holds the separator, a quote or a line break |
Transfer to a remote database
The pattern for a large load. Each step is a job: a transaction and a restart point.
- Stage. Load the file into a
temporarytable of the staging database withload!. The file is read once; a restart reads the stage, not the file; - Check and shape in SQL. Deduplicate, join reference data and total with SQL on SQLite;
- Prepare the target.
sales.prepare(tables)disables indexes, constraints and triggers; - Transfer in parallel groups, one group per level of the table hierarchy: head tables first, leaf tables last, in the order you write. Each aspect has its own connection and loads with
load!(mode: "upsert"), withcommit;and a checkpoint per batch; - Activate.
sales.activate(tables)andsales.reindex(tables)enable the constraints again. A constraint that fails names the rows that break it, and the data stays for inspection; - Record the counts of every step in the run log.
Upserts make every step safe to run twice: a failed run starts again from its last checkpoint and never duplicates a row.
Example: a file into the ERP
# load the orders of a shop file into the ERP
driver orders_upload is
class Order = {id: Integer, customer_id: Integer, amount: Decimal, created: Instant} <: Record is
method validate(@self) is
expect self.amount >= 0d;
return;
end Order;
table stage_orders: Order in "stage" is
key (id);
temporary;
end stage_orders;
table orders: Order in "erp" is ** "erp" = "eve://prod/erp" in eve.cfg
key (id);
end orders;
process main(path: String) is
stage: job
do
new bad := Table(:Reject)();
stage_orders.load!(csv.read!(:Order)(path, encoding: "utf-8", normalize: "nfc",
rejects: @bad, max_rejects: 100), mode: "upsert");
csv.write("out/orders_rejects.csv", bad) if bad.count() > 0;
done stage;
transfer: job
new since := checkpoint.read!(:Integer)("orders_upload", db: "erp", default: 0);
do
for batch in stage_orders.scan!("id > ? order by id", since).batch(5000) do
orders.load!(batch, mode: "upsert");
checkpoint.write("orders_upload", batch[-1].id, db: "erp");
commit; ** the batch and the checkpoint, together
done;
done transfer;
recover
retry after 5s if $error.transient and $error.attempt < 5;
abort;
return;
end orders_upload;
The stage job reads the file once. The checkpoint is written in the target database, so it commits with its batch; a failure of transfer restarts at the last committed batch.
Dat file; one rejects output per read or per run; NFC by default or off; the module of compression. The list is in plan/review_level5.md.