Skip to content
Open
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
57 changes: 25 additions & 32 deletions services/cubejs/src/routes/loadExport.js
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ import {
writeBinaryChunk,
writeRowStreamAsArrow,
} from "../utils/arrowSerializer.js";
import { execClickHouseArrowStream } from "../utils/clickhouseArrow.js";

const prepareAnnotation =
typeof prepareAnnotationModule.prepareAnnotation === "function"
Expand Down Expand Up @@ -54,20 +55,6 @@ function isClickHouseContext(securityContext) {
return typeof dbType === "string" && dbType.toLowerCase() === "clickhouse";
}

function removeTrailingSemicolon(query) {
const trimmed = String(query ?? "").trimEnd();
let lastNonSemiIdx = trimmed.length;
for (let i = lastNonSemiIdx; i > 0; i--) {
if (trimmed[i - 1] !== ";") {
lastNonSemiIdx = i;
break;
}
}
return lastNonSemiIdx !== trimmed.length
? trimmed.slice(0, lastNonSemiIdx)
: trimmed;
}

function normalizeClickHouseCSVLine(line) {
const normalized = line.replace(CLICKHOUSE_NULL_TOKEN_RE, "");
if (normalized.endsWith("\r\n")) return normalized;
Expand Down Expand Up @@ -426,13 +413,10 @@ async function executeNativeClickHouseCsv(
}

async function executeNativeClickHouseArrow(res, query, values, driver, signal) {
const result = await driver.client.exec({
query: `${removeTrailingSemicolon(sqlstring.format(query, values || []))}\nFORMAT ArrowStream`,
clickhouse_settings: {
...driver.config?.clickhouseSettings,
output_format_arrow_compression_method: "none",
},
abort_signal: signal,
const result = await execClickHouseArrowStream({
driver,
sql: sqlstring.format(query, values || []),
signal
});

const stream = typeof result.stream === "function"
Expand Down Expand Up @@ -496,17 +480,26 @@ async function tryHandleLoadExport(req, res, cubejs, query, format) {
const driver = await cubejs.options.driverFactory({
securityContext: plan.context.securityContext,
});
res.set(ARROW_HEADERS);
setNativeArrowFieldMappingHeaders(res, plan);
await executeNativeClickHouseArrow(
res,
nativeQuery.query,
nativeQuery.values,
driver,
abortController.signal
);
res.end();
return true;
try {
res.set(ARROW_HEADERS);
setNativeArrowFieldMappingHeaders(res, plan);
await executeNativeClickHouseArrow(
res,
nativeQuery.query,
nativeQuery.values,
driver,
abortController.signal
);
res.end();
return true;
} catch (err) {
if (abortController.signal.aborted) return true;
if (res.writableEnded) throw err;
console.warn(
"Native ClickHouse Arrow failed; falling back to semantic stream:",
err?.message || err
);
}
}

if (!canSemanticStreamLoadExport(format, plan.capabilities)) {
Expand Down
32 changes: 9 additions & 23 deletions services/cubejs/src/routes/runSql.js
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import crypto from "crypto";

import { loadRules } from "../utils/queryRewrite.js";
import { serializeRowsToArrow } from "../utils/arrowSerializer.js";
import { execClickHouseArrowStream } from "../utils/clickhouseArrow.js";
import { validateFormat } from "../utils/formatValidator.js";
import { writeRowsAsCSV, writeTextChunk } from "../utils/csvSerializer.js";
import { buildJSONStat } from "../utils/jsonstatBuilder.js";
Expand Down Expand Up @@ -62,20 +63,6 @@ function addUniqueColumns(target, names) {
}
}

function removeTrailingSemicolon(query) {
const trimmed = String(query ?? "").trimEnd();
let lastNonSemiIdx = trimmed.length;
for (let i = lastNonSemiIdx; i > 0; i--) {
if (trimmed[i - 1] !== ";") {
lastNonSemiIdx = i;
break;
}
}
return lastNonSemiIdx !== trimmed.length
? trimmed.slice(0, lastNonSemiIdx)
: trimmed;
}

function deriveExportColumnsFromRunSql(body, rows) {
if (rows.length > 0) {
return Object.keys(rows[0]);
Expand Down Expand Up @@ -213,20 +200,19 @@ export default async (req, res, cubejs) => {
}

if (format === "arrow" && isClickHouse(securityContext)) {
const clickhouseQuery = `${removeTrailingSemicolon(sql)}\nFORMAT ArrowStream`;
const result = await driver.client.exec({
query: clickhouseQuery,
clickhouse_settings: {
...driver.config?.clickhouseSettings,
output_format_arrow_compression_method: "none",
},
abort_signal: abortController.signal,
const result = await execClickHouseArrowStream({
driver,
sql,
signal: abortController.signal
});

res.set(ARROW_HEADERS);

try {
for await (const chunk of result.stream) {
const stream = typeof result.stream === "function"
? result.stream()
: result.stream;
for await (const chunk of stream) {
await writeBinaryChunk(res, chunk, abortController.signal);
}
} catch (streamErr) {
Expand Down
85 changes: 85 additions & 0 deletions services/cubejs/src/utils/__tests__/clickhouseArrow.test.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,85 @@
import { describe, it } from "node:test";
import assert from "node:assert/strict";

import {
execClickHouseArrowStream,
isForbiddenClickHouseArrowCompressionError
} from "../clickhouseArrow.js";

describe("clickhouseArrow", () => {
it("detects the readonly Arrow compression SET error", () => {
assert.equal(
isForbiddenClickHouseArrowCompressionError(
new Error(
"Cannot modify 'output_format_arrow_compression_method' setting in readonly mode."
)
),
true
);
assert.equal(
isForbiddenClickHouseArrowCompressionError(new Error("timeout")),
false
);
});

it("retries without the compression override after a readonly SET error", async () => {
const settingsByCall = [];
const driver = {
config: { clickhouseSettings: { max_execution_time: 30 } },
client: {
exec: async (opts) => {
settingsByCall.push(opts.clickhouse_settings);
if (settingsByCall.length === 1) {
throw new Error(
"Cannot modify 'output_format_arrow_compression_method' setting in readonly mode."
);
}
return { ok: true, query: opts.query };
}
}
};

const result = await execClickHouseArrowStream({
driver,
sql: "SELECT 1;",
signal: undefined
});

assert.equal(result.ok, true);
assert.match(result.query, /FORMAT ArrowStream/);
assert.equal(settingsByCall.length, 2);
assert.equal(settingsByCall[0].max_execution_time, 30);
assert.equal(settingsByCall[0].output_format_arrow_compression_method, "none");
assert.equal(settingsByCall[1].max_execution_time, 30);
assert.equal(
settingsByCall[1].output_format_arrow_compression_method,
undefined
);
});

it("skips the compression SET when the driver is readonly", async () => {
const settingsByCall = [];
const driver = {
config: { clickhouseSettings: {} },
readOnly: () => true,
client: {
exec: async (opts) => {
settingsByCall.push(opts.clickhouse_settings);
return { ok: true, query: opts.query };
}
}
};

await execClickHouseArrowStream({
driver,
sql: "SELECT 1",
signal: undefined
});

assert.equal(settingsByCall.length, 1);
assert.equal(
settingsByCall[0].output_format_arrow_compression_method,
undefined
);
});
});
66 changes: 66 additions & 0 deletions services/cubejs/src/utils/clickhouseArrow.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
function removeTrailingSemicolon(query) {
const trimmed = String(query ?? "").trimEnd();
let lastNonSemiIdx = trimmed.length;
for (let i = lastNonSemiIdx; i > 0; i--) {
if (trimmed[i - 1] !== ";") {
lastNonSemiIdx = i;
break;
}
}
return lastNonSemiIdx !== trimmed.length
? trimmed.slice(0, lastNonSemiIdx)
: trimmed;
}

function getClickHouseClient(driver) {
if (driver?.client && typeof driver.client.exec === "function") {
return driver.client;
}
if (driver && typeof driver.exec === "function") {
return driver;
}
return null;
}

export function isForbiddenClickHouseArrowCompressionError(err) {
const msg = String(err?.message || err);
return (
msg.includes("output_format_arrow_compression_method")
&& (msg.includes("readonly") || msg.includes("Cannot modify"))
);
}

/**
* Native ClickHouse ArrowStream. Prefer uncompressed IPC so the browser
* apache-arrow decoder can read it. Readonly ClickHouse users cannot SET
* that codec, even to the current value, so skip or retry without it.
*/
export async function execClickHouseArrowStream({ driver, sql, signal }) {
const client = getClickHouseClient(driver);
if (!client) {
throw new Error("ClickHouse driver has no exec client for ArrowStream");
}

const query = `${removeTrailingSemicolon(sql)}\nFORMAT ArrowStream`;
const baseSettings = { ...(driver.config?.clickhouseSettings || {}) };
const readonly = typeof driver.readOnly === "function" && driver.readOnly();

const run = (clickhouse_settings) => client.exec({
query,
clickhouse_settings,
abort_signal: signal
});

if (!readonly) {
try {
return await run({
...baseSettings,
output_format_arrow_compression_method: "none"
});
} catch (err) {
if (!isForbiddenClickHouseArrowCompressionError(err)) throw err;
}
}

return await run(baseSettings);
}
Loading