Skip to content

Commit 7d620a5

Browse files
committed
fix(clickhouse): classify streamed limit failures and keep canonical log fields
Spread caller-supplied logFields before the canonical ones so a caller cannot overwrite error, query, params, or queryId in a failure log. queryFastStream classifies quota failures the same way the buffered query paths do; a limit hit partway through a stream is still the caller asking for too much. The supplemental EXPLAIN queries carry the originating TSQL too.
1 parent 1b32811 commit 7d620a5

2 files changed

Lines changed: 10 additions & 3 deletions

File tree

internal-packages/clickhouse/src/client/client.ts

Lines changed: 9 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -332,12 +332,12 @@ export class ClickhouseClient implements ClickhouseReader, ClickhouseWriter {
332332

333333
if (clickhouseError) {
334334
const errorLogFields = {
335+
...req.logFields,
335336
name: req.name,
336337
error: clickhouseError,
337338
query: req.query,
338339
params,
339340
queryId,
340-
...req.logFields,
341341
};
342342

343343
if (isClickhouseQuotaError(clickhouseError)) {
@@ -623,13 +623,19 @@ export class ClickhouseClient implements ClickhouseReader, ClickhouseWriter {
623623

624624
span.setAttributes({ "clickhouse.rows": rowCount });
625625
} catch (error) {
626-
self.logger.error("Error streaming clickhouse", {
626+
const errorLogFields = {
627627
name: req.name,
628628
error,
629629
query: req.query,
630630
params,
631631
queryId,
632-
});
632+
};
633+
634+
if (error instanceof Error && isClickhouseQuotaError(error)) {
635+
self.logger.warn("Streamed query exceeded a ClickHouse limit", errorLogFields);
636+
} else {
637+
self.logger.error("Error streaming clickhouse", errorLogFields);
638+
}
633639

634640
if (error instanceof Error) {
635641
recordClickhouseError(span, error);

internal-packages/clickhouse/src/client/tsql.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -242,6 +242,7 @@ export async function executeTSQL<TOut extends z.ZodSchema>(
242242
params: z.record(z.any()),
243243
schema: z.object({ explain: z.string() }),
244244
settings: options.clickhouseSettings,
245+
logFields: { tsql: options.query },
245246
});
246247

247248
const [additionalError, additionalResult] = await additionalQueryFn(params);

0 commit comments

Comments
 (0)