Skip to content

Commit b50a518

Browse files
committed
feat: streaming CSV and Parquet export endpoints (#134)
- Add GET /transfers.csv and GET /transfers.parquet - Cursor-based batch fetching (500 rows/batch) keeps memory flat at any scale - CSV streams directly into response via @fast-csv/format - Parquet writes to tmp file then pipes to response, cleaned up after send - Both endpoints accept the same query params as the transfers route (address, contractId, fromLedger, toLedger, fromDate, toDate, eventType)
1 parent 1da3a82 commit b50a518

4 files changed

Lines changed: 239 additions & 0 deletions

File tree

package-lock.json

Lines changed: 56 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

package.json

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -60,13 +60,15 @@
6060
"@apollo/server": "^5.5.1",
6161
"@as-integrations/express4": "^1.1.2",
6262
"@asteasolutions/zod-to-openapi": "^8.5.0",
63+
"@fast-csv/format": "^5.0.7",
6364
"@prisma/client": "^5.10.0",
6465
"@stellar/stellar-sdk": "^15.0.1",
6566
"cors": "^2.8.5",
6667
"dotenv": "^16.4.5",
6768
"express": "^4.18.3",
6869
"express-rate-limit": "^8.3.2",
6970
"graphql": "^16.11.0",
71+
"parquetjs-lite": "^0.8.7",
7072
"ws": "^8.20.0",
7173
"zod": "^4.4.3"
7274
},

src/api.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ import { createAccountsRouter } from "./api/accounts";
99
import { createWebhooksRouter } from "./api/webhooks";
1010
import { createGraphQLMiddleware } from "./graphql/server";
1111
import { createPopularAssetsRouter } from "./routes/assets/popular";
12+
import { createExportsRouter } from "./routes/exports";
1213
import {
1314
hostFnQuerySchema,
1415
nftOwnerParamsSchema,
@@ -96,6 +97,9 @@ export function createApp(): express.Application {
9697
// ── Assets routes ───────────────────────────────────────────────────────────
9798
app.use("/assets", createPopularAssetsRouter());
9899

100+
// ── Export routes ─────────────────────────────────────────────────────────────
101+
app.use("/", createExportsRouter());
102+
99103
// ── Helpers ──────────────────────────────────────────────────────────────────
100104
const parseIntParam = (val: unknown, fallback: number): number => {
101105
const n = parseInt(String(val), 10);

src/routes/exports.ts

Lines changed: 177 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,177 @@
1+
import { Router, Request, Response, NextFunction } from "express";
2+
import { format as csvFormat } from "@fast-csv/format";
3+
import { prisma, toDisplayAmount } from "../db";
4+
import os from "os";
5+
import path from "path";
6+
import fs from "fs";
7+
8+
// How many rows we fetch per DB round-trip. Keeps memory flat.
9+
const BATCH_SIZE = 500;
10+
11+
// ── Shared: parse query params into a Prisma where clause ────────────────────
12+
function buildWhere(query: Record<string, unknown>) {
13+
const {
14+
address,
15+
contractId,
16+
fromLedger,
17+
toLedger,
18+
fromDate,
19+
toDate,
20+
eventType,
21+
} = query;
22+
23+
const where: Record<string, unknown> = {};
24+
25+
if (address) {
26+
where.OR = [{ fromAddress: address }, { toAddress: address }];
27+
}
28+
if (contractId) where.contractId = contractId;
29+
if (eventType) {
30+
const types = String(eventType).split(",").map((s) => s.trim()).filter(Boolean);
31+
if (types.length) where.eventType = { in: types };
32+
}
33+
34+
const ledgerRange: Record<string, number> = {};
35+
if (fromLedger) ledgerRange.gte = parseInt(String(fromLedger), 10);
36+
if (toLedger) ledgerRange.lte = parseInt(String(toLedger), 10);
37+
if (Object.keys(ledgerRange).length) where.ledger = ledgerRange;
38+
39+
const dateRange: Record<string, Date> = {};
40+
if (fromDate) dateRange.gte = new Date(String(fromDate));
41+
if (toDate) dateRange.lte = new Date(String(toDate));
42+
if (Object.keys(dateRange).length) where.ledgerClosedAt = dateRange;
43+
44+
return where;
45+
}
46+
47+
// ── Shared: async generator that yields rows in batches via cursor ────────────
48+
async function* streamTransfers(where: Record<string, unknown>) {
49+
let lastId: number | undefined = undefined;
50+
51+
while (true) {
52+
const rows: Awaited<ReturnType<typeof prisma.tokenTransfer.findMany>> = await prisma.tokenTransfer.findMany({
53+
where,
54+
orderBy: { id: "asc" },
55+
take: BATCH_SIZE,
56+
...(lastId !== undefined ? { cursor: { id: lastId }, skip: 1 } : {}),
57+
});
58+
59+
if (rows.length === 0) break;
60+
61+
for (const row of rows) {
62+
yield row;
63+
}
64+
65+
if (rows.length < BATCH_SIZE) break;
66+
lastId = rows[rows.length - 1].id;
67+
}
68+
}
69+
70+
// ── CSV endpoint ─────────────────────────────────────────────────────────────
71+
async function handleCsvExport(req: Request, res: Response, next: NextFunction) {
72+
try {
73+
const where = buildWhere(req.query as Record<string, unknown>);
74+
75+
res.setHeader("Content-Type", "text/csv");
76+
res.setHeader("Content-Disposition", "attachment; filename=\"transfers.csv\"");
77+
res.setHeader("Transfer-Encoding", "chunked");
78+
79+
const csvStream = csvFormat({ headers: true });
80+
csvStream.pipe(res);
81+
82+
for await (const row of streamTransfers(where)) {
83+
csvStream.write({
84+
id: row.id,
85+
contractId: row.contractId,
86+
eventType: row.eventType,
87+
fromAddress: row.fromAddress ?? "",
88+
toAddress: row.toAddress ?? "",
89+
amount: row.amount,
90+
displayAmount: toDisplayAmount(row.amount),
91+
ledger: row.ledger,
92+
ledgerClosedAt: row.ledgerClosedAt.toISOString(),
93+
txHash: row.txHash,
94+
eventId: row.eventId,
95+
isSac: row.isSac ?? false,
96+
createdAt: row.createdAt.toISOString(),
97+
});
98+
}
99+
100+
csvStream.end();
101+
} catch (err) {
102+
next(err);
103+
}
104+
}
105+
106+
// ── Parquet endpoint ─────────────────────────────────────────────────────────
107+
async function handleParquetExport(req: Request, res: Response, next: NextFunction) {
108+
// parquetjs-lite is a CommonJS module — require() avoids ESM interop issues
109+
// eslint-disable-next-line @typescript-eslint/no-var-requires
110+
const parquet = require("parquetjs-lite");
111+
112+
const tmpFile = path.join(os.tmpdir(), `transfers-${Date.now()}-${Math.random().toString(36).slice(2)}.parquet`);
113+
114+
try {
115+
const where = buildWhere(req.query as Record<string, unknown>);
116+
117+
const schema = new parquet.ParquetSchema({
118+
id: { type: "INT64" },
119+
contractId: { type: "UTF8" },
120+
eventType: { type: "UTF8" },
121+
fromAddress: { type: "UTF8", optional: true },
122+
toAddress: { type: "UTF8", optional: true },
123+
amount: { type: "UTF8" },
124+
displayAmount: { type: "UTF8" },
125+
ledger: { type: "INT32" },
126+
ledgerClosedAt: { type: "UTF8" },
127+
txHash: { type: "UTF8" },
128+
eventId: { type: "UTF8" },
129+
isSac: { type: "BOOLEAN", optional: true },
130+
createdAt: { type: "UTF8" },
131+
});
132+
133+
const writer = await parquet.ParquetWriter.openFile(schema, tmpFile);
134+
135+
for await (const row of streamTransfers(where)) {
136+
await writer.appendRow({
137+
id: row.id,
138+
contractId: row.contractId,
139+
eventType: row.eventType,
140+
fromAddress: row.fromAddress ?? null,
141+
toAddress: row.toAddress ?? null,
142+
amount: row.amount,
143+
displayAmount: toDisplayAmount(row.amount),
144+
ledger: row.ledger,
145+
ledgerClosedAt: row.ledgerClosedAt.toISOString(),
146+
txHash: row.txHash,
147+
eventId: row.eventId,
148+
isSac: row.isSac ?? null,
149+
createdAt: row.createdAt.toISOString(),
150+
});
151+
}
152+
153+
await writer.close();
154+
155+
res.setHeader("Content-Type", "application/octet-stream");
156+
res.setHeader("Content-Disposition", "attachment; filename=\"transfers.parquet\"");
157+
158+
const fileStream = fs.createReadStream(tmpFile);
159+
fileStream.pipe(res);
160+
fileStream.on("end", () => fs.unlink(tmpFile, () => {}));
161+
fileStream.on("error", (err) => {
162+
fs.unlink(tmpFile, () => {});
163+
next(err);
164+
});
165+
} catch (err) {
166+
fs.unlink(tmpFile, () => {});
167+
next(err);
168+
}
169+
}
170+
171+
// ── Router ───────────────────────────────────────────────────────────────────
172+
export function createExportsRouter(): Router {
173+
const router = Router();
174+
router.get("/transfers.csv", handleCsvExport);
175+
router.get("/transfers.parquet", handleParquetExport);
176+
return router;
177+
}

0 commit comments

Comments
 (0)