mirror of
https://github.com/R0m1k3/CollectFlow.git
synced 2026-10-11 17:26:32 +02:00
fix(api): tient les très gros fournisseurs (132k articles)
Le journal du fournisseur D005 montre 132 388 lignes persistées pour un seul
fournisseur. À cette échelle, /api/v1/grid sans limit ne tenait pas.
Deux plantages francs, indépendants de la taille de la réponse :
- getNetworkMetricsByCodeCentrale : inArray avec 132 000 valeurs, soit autant
de paramètres liés, alors que le protocole PostgreSQL en accepte 65535.
- pgGetGammesByCodeins : même problème sur son IN (…).
Les deux lectures sont désormais découpées par lots de 2000. Le bug existait
pour tout appel dépassant ~65 000 codes, indépendamment de l'API.
Volume : charger 132k lignes complètes puis les sérialiser d'un bloc représente
plusieurs centaines de Mo en mémoire. Sans `limit`, la réponse est maintenant
diffusée en flux, par lots internes de 1000, à mémoire constante. La forme
{ data, pagination, meta } est inchangée — un client existant ne voit que des
octets arrivant progressivement.
Une erreur survenant après les premiers octets ne peut plus changer le code
HTTP : le document est alors fermé avec meta.erreur et meta.complet=false, pour
qu'une réponse tronquée ne passe pas pour complète.
`fields` reste le levier décisif à cette échelle : une ligne complète porte les
séries mensuelles et les ventilations par magasin.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Y26nRZxTR57K7h8yqsF675
This commit is contained in:
3 files changed
+135
-26
No files matched your search
@@ -77,7 +77,7 @@ export async function GET(req: NextRequest) {
|
||||
computedOnDemand = true;
|
||||
}
|
||||
|
||||
const result = await queryGridRows({
|
||||
const filtres = {
|
||||
codeFournisseur: q.fournisseur,
|
||||
gamme: q.gamme,
|
||||
code1: q.code1,
|
||||
@@ -86,16 +86,21 @@ export async function GET(req: NextRequest) {
|
||||
search: q.search,
|
||||
sort: toSortKey(q.sort),
|
||||
order: q.order,
|
||||
page: q.page,
|
||||
limit: q.limit,
|
||||
});
|
||||
};
|
||||
|
||||
// Sans `limit`, la réponse peut porter tout un fournisseur — et certains en
|
||||
// comptent plus de 130 000. Tout charger puis tout sérialiser d'un bloc ferait
|
||||
// exploser la mémoire du processus : on diffuse par lots, à mémoire constante.
|
||||
if (q.limit == null) {
|
||||
return streamAllRows(req, filtres, q, freshness, computedOnDemand);
|
||||
}
|
||||
|
||||
const result = await queryGridRows({ ...filtres, page: q.page, limit: q.limit });
|
||||
|
||||
// Métriques Qlik et gamme serveur relues au moment de l'appel (voir api-enrich).
|
||||
const rows = q.enrich === "1" ? await enrichRows(result.rows) : result.rows;
|
||||
|
||||
// Sans `limit`, tout a été renvoyé : la pagination doit le refléter (une seule
|
||||
// page couvrant le total), sinon `hasMore` mentirait.
|
||||
const pagination = buildPagination(q.page, q.limit ?? Math.max(1, result.total), result.total);
|
||||
const pagination = buildPagination(q.page, q.limit, result.total);
|
||||
|
||||
return ok(
|
||||
rows.map((r) => pickFields(r, q.fields)),
|
||||
@@ -127,3 +132,91 @@ export async function GET(req: NextRequest) {
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
/** Lots internes : compromis entre nombre d'allers-retours SQL et mémoire retenue. */
|
||||
const STREAM_BATCH = 1000;
|
||||
|
||||
type Filtres = Parameters<typeof queryGridRows>[0];
|
||||
|
||||
/**
|
||||
* Diffuse **toutes** les lignes du fournisseur en une seule réponse JSON, sans
|
||||
* jamais la construire entièrement en mémoire.
|
||||
*
|
||||
* La forme reste celle des autres réponses (`{ data, pagination, meta }`) : un
|
||||
* client existant ne voit aucune différence, il reçoit simplement les octets au
|
||||
* fil de l'eau. `pagination` et `meta` sont écrits en dernier, une fois le total
|
||||
* connu — l'ordre des clés n'a aucune importance en JSON.
|
||||
*
|
||||
* Une erreur survenant après les premiers octets ne peut plus changer le code
|
||||
* HTTP : on ferme alors le document avec `meta.erreur`, pour que l'appelant
|
||||
* détecte une réponse tronquée au lieu de la croire complète.
|
||||
*/
|
||||
function streamAllRows(
|
||||
req: NextRequest,
|
||||
filtres: Filtres,
|
||||
q: { enrich: string; fields?: string; fournisseur: string },
|
||||
freshness: { rowCount: number; computedAt: string | null },
|
||||
computedOnDemand: boolean,
|
||||
): Response {
|
||||
const encoder = new TextEncoder();
|
||||
|
||||
const stream = new ReadableStream<Uint8Array>({
|
||||
async start(controller) {
|
||||
const write = (s: string) => controller.enqueue(encoder.encode(s));
|
||||
let envoyees = 0;
|
||||
let total = 0;
|
||||
let computedAt: string | null = null;
|
||||
let erreur: string | null = null;
|
||||
|
||||
write('{"data":[');
|
||||
try {
|
||||
for (let page = 1; ; page++) {
|
||||
const lot = await queryGridRows({ ...filtres, page, limit: STREAM_BATCH });
|
||||
if (page === 1) {
|
||||
total = lot.total;
|
||||
computedAt = lot.computedAt;
|
||||
}
|
||||
if (lot.rows.length === 0) break;
|
||||
|
||||
const enrichies = q.enrich === "1" ? await enrichRows(lot.rows) : lot.rows;
|
||||
for (const r of enrichies) {
|
||||
write((envoyees === 0 ? "" : ",") + JSON.stringify(pickFields(r, q.fields)));
|
||||
envoyees++;
|
||||
}
|
||||
if (lot.rows.length < STREAM_BATCH) break;
|
||||
}
|
||||
} catch (e) {
|
||||
erreur = e instanceof Error ? e.message : String(e);
|
||||
console.error(`[api/v1/grid] diffusion interrompue pour ${q.fournisseur} après ${envoyees} lignes:`, erreur);
|
||||
}
|
||||
|
||||
const complet = erreur === null && envoyees === total;
|
||||
write("]," + JSON.stringify({
|
||||
pagination: { page: 1, limit: envoyees, total, totalPages: 1, hasMore: !complet },
|
||||
meta: {
|
||||
fournisseur: q.fournisseur,
|
||||
computedAt,
|
||||
snapshotComputedAt: freshness.computedAt,
|
||||
snapshotRowCount: freshness.rowCount,
|
||||
enrichi: q.enrich === "1",
|
||||
computedOnDemand,
|
||||
diffuse: true,
|
||||
complet,
|
||||
...(erreur
|
||||
? { erreur: `Diffusion interrompue après ${envoyees} lignes sur ${total} : ${erreur}` }
|
||||
: {}),
|
||||
},
|
||||
}).slice(1));
|
||||
controller.close();
|
||||
},
|
||||
});
|
||||
|
||||
return new Response(stream, {
|
||||
headers: {
|
||||
"Content-Type": "application/json; charset=utf-8",
|
||||
"Cache-Control": "no-store",
|
||||
// Empêche un proxy de tamponner la réponse entière avant de la relayer.
|
||||
"X-Accel-Buffering": "no",
|
||||
},
|
||||
});
|
||||
}
|
||||
+22
-15
@@ -314,22 +314,29 @@ export async function pgGetGammesByCodeins(codeins: string[]): Promise<Map<strin
|
||||
const unique = [...new Set(codeins.map(c => String(c ?? "").trim()).filter(Boolean))];
|
||||
if (unique.length === 0) return new Map();
|
||||
|
||||
// DISTINCT ON (codein) → 1 gamme par article, saison la plus récente en premier.
|
||||
const result = await pgNoParallel(sql`
|
||||
SELECT DISTINCT ON (a.codein)
|
||||
TRIM(a.codein::text) AS codein,
|
||||
g.code AS gamme_code
|
||||
FROM art_gamme_saison ags
|
||||
JOIN articles a ON a.no_id = ags.artnoid
|
||||
JOIN gammes g ON g.no_id = ags.idgamme
|
||||
JOIN saisons s ON s.no_id = ags.idsaison
|
||||
WHERE TRIM(a.codein::text) IN (${sql.join(unique.map(c => sql`${c}`), sql`, `)})
|
||||
ORDER BY a.codein, s.no_id DESC
|
||||
`);
|
||||
|
||||
// Découpage obligatoire : chaque codein devient un paramètre lié, or PostgreSQL
|
||||
// en accepte 65535 au maximum. Un gros fournisseur dépasse 130 000 articles —
|
||||
// la requête échouait alors purement et simplement.
|
||||
const CHUNK = 2000;
|
||||
const map = new Map<string, string>();
|
||||
for (const row of result.rows as unknown as { codein: string; gamme_code: string }[]) {
|
||||
if (row.codein && row.gamme_code) map.set(row.codein, String(row.gamme_code).trim());
|
||||
|
||||
for (let i = 0; i < unique.length; i += CHUNK) {
|
||||
const batch = unique.slice(i, i + CHUNK);
|
||||
// DISTINCT ON (codein) → 1 gamme par article, saison la plus récente en premier.
|
||||
const result = await pgNoParallel(sql`
|
||||
SELECT DISTINCT ON (a.codein)
|
||||
TRIM(a.codein::text) AS codein,
|
||||
g.code AS gamme_code
|
||||
FROM art_gamme_saison ags
|
||||
JOIN articles a ON a.no_id = ags.artnoid
|
||||
JOIN gammes g ON g.no_id = ags.idgamme
|
||||
JOIN saisons s ON s.no_id = ags.idsaison
|
||||
WHERE TRIM(a.codein::text) IN (${sql.join(batch.map(c => sql`${c}`), sql`, `)})
|
||||
ORDER BY a.codein, s.no_id DESC
|
||||
`);
|
||||
for (const row of result.rows as unknown as { codein: string; gamme_code: string }[]) {
|
||||
if (row.codein && row.gamme_code) map.set(row.codein, String(row.gamme_code).trim());
|
||||
}
|
||||
}
|
||||
console.log(`[pg-ff] pgGetGammesByCodeins: ${map.size}/${unique.length} articles avec gamme`);
|
||||
return map;
|
||||
|
||||
@@ -42,10 +42,19 @@ export async function getNetworkMetricsByCodeCentrale(
|
||||
const unique = [...new Set(codes.filter(Boolean))];
|
||||
if (unique.length === 0) return out;
|
||||
|
||||
const rows = await db
|
||||
.select()
|
||||
.from(qlikNetworkMetrics)
|
||||
.where(inArray(qlikNetworkMetrics.codeCentrale, unique));
|
||||
// Découpage obligatoire : `inArray` produit un paramètre lié par valeur, or le
|
||||
// protocole PostgreSQL en accepte 65535 au maximum. Un gros fournisseur peut
|
||||
// dépasser 130 000 codes — la requête échouait alors purement et simplement.
|
||||
const CHUNK = 2000;
|
||||
const rows: Array<typeof qlikNetworkMetrics.$inferSelect> = [];
|
||||
for (let i = 0; i < unique.length; i += CHUNK) {
|
||||
const batch = unique.slice(i, i + CHUNK);
|
||||
const part = await db
|
||||
.select()
|
||||
.from(qlikNetworkMetrics)
|
||||
.where(inArray(qlikNetworkMetrics.codeCentrale, batch));
|
||||
rows.push(...part);
|
||||
}
|
||||
|
||||
for (const r of rows) {
|
||||
out.set(r.codeCentrale, {
|
||||
|
||||
Reference in new issue
Block a user