Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 18 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -139,6 +139,24 @@ jobs:
print("round-trip figures preserved:", nav(r), npos(r), spct(r))
PY

- name: Large-file - streaming aggregate / split / delta (constant memory)
run: |
set -e
python3 -m pip install --quiet lxml
P=Large_File_Processing/python
python3 $P/make_large_sample.py big.xml 30000
xmllint --noout --nonet --schema .schema-cache/4.2.9/FundsXML.xsd big.xml
python3 $P/stream_aggregate.py big.xml | tee agg.txt
grep -q '^positions : 30000$' agg.txt
grep -q '^sum value (EUR): 30000000.00$' agg.txt
javac -d /tmp/lf Large_File_Processing/java/StreamAggregate.java
java -Xmx64m -cp /tmp/lf StreamAggregate big.xml | tee aggj.txt
grep -q '^positions : 30000$' aggj.txt
python3 $P/split.py big.xml chunks/ 10000
test "$(ls chunks/chunk-*.xml | wc -l)" = "3"
xmllint --noout --nonet --schema .schema-cache/4.2.9/FundsXML.xsd chunks/chunk-0001.xml
python3 $P/delta_diff.py big.xml big.xml # identical -> exit 0

- name: Regression - legacy XSLT 1.0 report still runs
run: |
xsltproc XSLT_DataQuality_Checks/Enhanced_Check/FundsXML_CompleteDQReport_HTML.xsl \
Expand Down
4 changes: 4 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -24,5 +24,9 @@ out_*.csv
out_*.xml
regen.xml
regenerated.xml
big.xml
chunks/
agg.txt
aggj.txt
*.db
*.svrl
54 changes: 54 additions & 0 deletions Large_File_Processing/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
# Large-File / Streaming Processing

![status](https://img.shields.io/badge/Python%20%2B%20Java%20StAX-verified-brightgreen)

Enterprise FundsXML feeds can be hundreds of MB. Loading them into a DOM blows
memory; these examples process them **streaming, at constant memory** —
verified: ~16 MiB RSS (Python) / ~2 MiB heap (Java, `-Xmx64m`) for a 20 000-
position file, independent of file size.

| Tool | What | Status |
|------|------|--------|
| [`python/make_large_sample.py`](python/make_large_sample.py) | Generate a big XSD-valid file by streaming **writes** | ✅ |
| [`python/stream_aggregate.py`](python/stream_aggregate.py) | `lxml.iterparse` aggregation (clears parsed siblings) | ✅ |
| [`java/StreamAggregate.java`](java/StreamAggregate.java) | StAX pull parser, native Java (no JAXB/DOM) | ✅ |
| [`python/split.py`](python/split.py) | Split into independently **XSD-valid** chunks | ✅ |
| [`python/delta_diff.py`](python/delta_diff.py) | INITIAL-vs-DELTA position diff (added/removed/changed) | ✅ |

FundsXML 4.x has no XML namespace — matchers use the bare `Position` tag.
Parsers disable DTD/external entities (XXE-safe).

## Run

```bash
# 1. synthetic 50k-position file (~constant-memory writer)
python3 Large_File_Processing/python/make_large_sample.py big.xml 50000

# 2. aggregate — Python and Java give identical totals at flat memory
python3 Large_File_Processing/python/stream_aggregate.py big.xml
javac -d /tmp/lf Large_File_Processing/java/StreamAggregate.java
java -Xmx64m -cp /tmp/lf StreamAggregate big.xml

# 3. split into XSD-valid chunks of 10k positions
python3 Large_File_Processing/python/split.py big.xml chunks/ 10000
xmllint --noout --schema .schema-cache/4.2.9/FundsXML.xsd chunks/chunk-0001.xml

# 4. day-over-day position delta (exit 1 if anything changed)
python3 Large_File_Processing/python/delta_diff.py yesterday.xml today.xml
```

## ETL note

The streaming readers are the natural building block for ETL pipelines
(Apache Camel / NiFi / Spring Batch): a `split` step feeds a parallel
`stream_aggregate`/load stage, and `delta_diff` drives incremental upserts when
a feed alternates `DataOperation` INITIAL/DELTA. Example Camel route shape:

```
from("file:in?include=.*\\.xml")
.to("exec:python3?args=Large_File_Processing/python/split.py ${file} work/")
.split(...).parallelProcessing()
.to("exec:...stream_aggregate.py ...") // or a JDBC load (see Database_Integration/)
```

CI runs steps 2–4 on a generated file (totals asserted, a chunk XSD-validated).
72 changes: 72 additions & 0 deletions Large_File_Processing/java/StreamAggregate.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
// Constant-memory aggregation over a huge FundsXML file via StAX (native Java,
// no JAXB, no DOM). Pull parser — never holds more than the current element.
//
// javac -d /tmp/lf Large_File_Processing/java/StreamAggregate.java
// java -cp /tmp/lf StreamAggregate <fundsxml.xml>
//
// FundsXML 4.x has no XML namespace — match the bare "Position" local name.
// Security: external entities / DTDs disabled (XXE-safe).

import java.io.FileInputStream;
import javax.xml.stream.XMLInputFactory;
import javax.xml.stream.XMLStreamConstants;
import javax.xml.stream.XMLStreamReader;

public class StreamAggregate {
public static void main(String[] args) throws Exception {
if (args.length != 1) {
System.err.println("usage: StreamAggregate <fundsxml.xml>");
System.exit(2);
}

XMLInputFactory xif = XMLInputFactory.newInstance();
xif.setProperty(XMLInputFactory.IS_SUPPORTING_EXTERNAL_ENTITIES, false);
xif.setProperty(XMLInputFactory.SUPPORT_DTD, false);

long n = 0;
double sumValue = 0, sumPct = 0;
boolean inPosition = false, inTotalValue = false;
String currentLeaf = null;

try (FileInputStream in = new FileInputStream(args[0])) {
XMLStreamReader r = xif.createXMLStreamReader(in);
while (r.hasNext()) {
switch (r.next()) {
case XMLStreamConstants.START_ELEMENT: {
String name = r.getLocalName();
if ("Position".equals(name)) { inPosition = true; n++; }
else if ("TotalValue".equals(name)) inTotalValue = true;
currentLeaf = name;
break;
}
case XMLStreamConstants.CHARACTERS: {
if (!inPosition || r.isWhiteSpace()) break;
String t = r.getText().trim();
if (t.isEmpty()) break;
if (inTotalValue && "Amount".equals(currentLeaf))
sumValue += Double.parseDouble(t);
else if ("TotalPercentage".equals(currentLeaf))
sumPct += Double.parseDouble(t);
break;
}
case XMLStreamConstants.END_ELEMENT: {
String name = r.getLocalName();
if ("Position".equals(name)) inPosition = false;
else if ("TotalValue".equals(name)) inTotalValue = false;
currentLeaf = null;
break;
}
default:
}
}
r.close();
}

Runtime rt = Runtime.getRuntime();
long usedMiB = (rt.totalMemory() - rt.freeMemory()) / (1024 * 1024);
System.out.printf("positions : %d%n", n);
System.out.printf("sum value (EUR): %.2f%n", sumValue);
System.out.printf("sum percentage : %.4f%n", sumPct);
System.out.printf("heap in use : %d MiB%n", usedMiB);
}
}
87 changes: 87 additions & 0 deletions Large_File_Processing/python/delta_diff.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
#!/usr/bin/env python3
"""Position-level diff between two FundsXML files (INITIAL vs DELTA / day-over-day).

Usage: delta_diff.py <old.xml> <new.xml> [--json]

FundsXML's ControlData/DataOperation says INITIAL (full snapshot) or DELTA
(changes only). Independent of that flag, this compares the two position sets
by UniqueID and reports:
added present in new, not in old
removed present in old, not in new
changed value or percentage differs (> 0.005 tolerance)
unchanged count only

Both files are streamed with iterparse (constant memory). Exit code: 0 if the
sets are identical, 1 if there are differences, 2 on usage error.
"""
import json
import sys

from lxml import etree


def index(path):
op = None
out = {}
ctx = etree.iterparse(path, events=("end",), resolve_entities=False,
no_network=True)
for _, el in ctx:
if el.tag == "DataOperation":
op = (el.text or "").strip()
elif el.tag == "Position":
uid = el.findtext("UniqueID")
amt = el.find("./TotalValue/Amount")
val = float(amt.text) if amt is not None and amt.text else 0.0
pct = float(el.findtext("TotalPercentage") or 0.0)
if uid:
out[uid] = (val, pct)
el.clear()
while el.getprevious() is not None:
del el.getparent()[0]
return op, out


def main() -> int:
args = [a for a in sys.argv[1:] if a != "--json"]
as_json = "--json" in sys.argv
if len(args) != 2:
print("usage: delta_diff.py <old.xml> <new.xml> [--json]",
file=sys.stderr)
return 2

op_old, a = index(args[0])
op_new, b = index(args[1])

added = sorted(set(b) - set(a))
removed = sorted(set(a) - set(b))
changed = sorted(u for u in (set(a) & set(b))
if abs(a[u][0] - b[u][0]) > 0.005
or abs(a[u][1] - b[u][1]) > 0.005)
unchanged = len(set(a) & set(b)) - len(changed)

if as_json:
print(json.dumps({
"old": {"file": args[0], "dataOperation": op_old, "positions": len(a)},
"new": {"file": args[1], "dataOperation": op_new, "positions": len(b)},
"added": added, "removed": removed,
"changed": [{"id": u, "old": a[u], "new": b[u]} for u in changed],
"unchanged": unchanged,
}, indent=2))
else:
print(f"old: {args[0]} DataOperation={op_old} positions={len(a)}")
print(f"new: {args[1]} DataOperation={op_new} positions={len(b)}")
print(f"added : {len(added)} {added[:10]}{' ...' if len(added) > 10 else ''}")
print(f"removed : {len(removed)} {removed[:10]}{' ...' if len(removed) > 10 else ''}")
print(f"changed : {len(changed)}")
for u in changed[:10]:
print(f" {u}: value {a[u][0]:.2f} -> {b[u][0]:.2f}, "
f"pct {a[u][1]:.4f} -> {b[u][1]:.4f}")
if len(changed) > 10:
print(" ...")
print(f"unchanged: {unchanged}")

return 0 if not (added or removed or changed) else 1


if __name__ == "__main__":
sys.exit(main())
83 changes: 83 additions & 0 deletions Large_File_Processing/python/make_large_sample.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
#!/usr/bin/env python3
"""Generate a large but XSD-valid FundsXML positions file by STREAMING writes.

Usage: make_large_sample.py <out.xml> [n_positions=50000]

Builds the document with plain incremental writes (no in-memory DOM), so it
scales to very large files at constant memory — the write-side counterpart to
the streaming readers in this directory. Output validates against the official
4.2.9 schema (Fund + Portfolio/Positions, each an Equity with Units).
"""
import sys

SCHEMA = ("https://github.com/fundsxml/schema/releases/download/"
"4.2.9/FundsXML.xsd")


def main() -> int:
if len(sys.argv) < 2:
print("usage: make_large_sample.py <out.xml> [n]", file=sys.stderr)
return 2
out = sys.argv[1]
n = int(sys.argv[2]) if len(sys.argv) > 2 else 50000

# Integer cents keep the value sum exact; percentages are derived so they
# sum to exactly 100 (last position absorbs the rounding remainder).
unit_value = 1000 # EUR per position
total = n * unit_value
base_pct = round(100.0 / n, 6)

with open(out, "w", encoding="utf-8") as f:
w = f.write
w('<?xml version="1.0" encoding="UTF-8"?>\n')
w('<FundsXML4 xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"\n')
w(f' xsi:noNamespaceSchemaLocation="{SCHEMA}">\n')
w(' <ControlData>\n')
w(f' <UniqueDocumentID>FUNDSXML_LARGE_{n}</UniqueDocumentID>\n')
w(' <DocumentGenerated>2025-10-02T00:00:00</DocumentGenerated>\n')
w(' <Version>4.2.9</Version>\n')
w(' <ContentDate>2025-10-01</ContentDate>\n')
w(' <DataSupplier><SystemCountry>AT</SystemCountry>'
'<Short>EURAM</Short><Name>Erste Asset Management GmbH</Name>'
'<Type>Asset Manager</Type></DataSupplier>\n')
w(' <DataOperation>INITIAL</DataOperation>\n')
w(' </ControlData>\n')
w(' <Funds><Fund>\n')
w(' <Identifiers><LEI>529900T8BM49AURSDO55</LEI></Identifiers>\n')
w(' <Names><OfficialName>Erste Large Synthetic Fund</OfficialName></Names>\n')
w(' <Currency>EUR</Currency>\n')
w(' <SingleFundFlag>true</SingleFundFlag>\n')
w(' <FundDynamicData>\n')
w(' <TotalAssetValues><TotalAssetValue>\n')
w(' <NavDate>2025-10-01</NavDate>\n')
w(' <TotalAssetNature>OFFICIAL</TotalAssetNature>\n')
w(f' <TotalNetAssetValue><Amount ccy="EUR">{total}.00</Amount>'
'</TotalNetAssetValue>\n')
w(' </TotalAssetValue></TotalAssetValues>\n')
w(' <Portfolios><Portfolio>\n')
w(' <NavDate>2025-10-01</NavDate>\n')
w(' <Positions>\n')
acc = 0.0
for i in range(1, n + 1):
if i < n:
pct = base_pct
acc += pct
else:
pct = round(100.0 - acc, 6) # remainder on the last one
w(f'<Position><UniqueID>ID_{i:08d}</UniqueID>'
f'<Currency>EUR</Currency>'
f'<TotalValue><Amount ccy="EUR">{unit_value}.00</Amount></TotalValue>'
f'<TotalPercentage>{pct:.6f}</TotalPercentage>'
f'<Equity><Units>{unit_value}.00</Units></Equity></Position>\n')
w(' </Positions>\n')
w(' </Portfolio></Portfolios>\n')
w(' </FundDynamicData>\n')
w(' </Fund></Funds>\n')
w('</FundsXML4>\n')

print(f"wrote {out}: {n} positions, total NAV {total}.00 EUR")
return 0


if __name__ == "__main__":
sys.exit(main())
Loading
Loading