From ee2663d5b287f5ebd6f80bc4f57c2fecba0c2eb3 Mon Sep 17 00:00:00 2001 From: 28winz-bot <28.winz@gmail.com> Date: Sat, 4 Jul 2026 10:36:46 +0700 Subject: [PATCH] fix(clickhouse-io): use official @clickhouse/client instead of unmaintained clickhouse package (#2422) * fix(clickhouse-io): use official @clickhouse/client instead of unmaintained clickhouse package The example imported the third-party `clickhouse` (TimonKK) package and used its API (new ClickHouse, .query().toPromise(), .insert().stream()). Migrate to the official @clickhouse/client: createClient() and structured clickhouse.insert({ table, values, format }). The structured values array also removes the previous SQL string-interpolation anti-pattern. Applies to the source skill and the ja-JP, zh-CN, zh-TW, ko-KR translated copies (translated code comments preserved). Refs: https://clickhouse.com/docs/integrations/javascript * fix(clickhouse-io): migrate remaining insert calls to @clickhouse/client Address review feedback: the earlier commit missed two spots that still used the legacy clickhouse API. - CDC example: clickhouse.insert('market_updates', [...]) -> clickhouse.insert({ table, values, format: 'JSONEachRow' }) - Single-row insertTrade: map the row to the column shape (same as the bulk path) instead of passing the raw trade object. Applies to the source skill and the ja-JP, zh-CN, zh-TW, ko-KR copies. --- docs/ja-JP/skills/clickhouse-io/SKILL.md | 87 +++++++++++++----------- docs/ko-KR/skills/clickhouse-io/SKILL.md | 84 ++++++++++++----------- docs/zh-CN/skills/clickhouse-io/SKILL.md | 87 +++++++++++++----------- docs/zh-TW/skills/clickhouse-io/SKILL.md | 87 +++++++++++++----------- skills/clickhouse-io/SKILL.md | 87 +++++++++++++----------- 5 files changed, 230 insertions(+), 202 deletions(-) diff --git a/docs/ja-JP/skills/clickhouse-io/SKILL.md b/docs/ja-JP/skills/clickhouse-io/SKILL.md index 59f071390..b1f9bdae6 100644 --- a/docs/ja-JP/skills/clickhouse-io/SKILL.md +++ b/docs/ja-JP/skills/clickhouse-io/SKILL.md @@ -151,39 +151,43 @@ ORDER BY market_id, date; ### 一括挿入(推奨) ```typescript -import { ClickHouse } from 'clickhouse' +import { createClient } from '@clickhouse/client' -const clickhouse = new ClickHouse({ - url: process.env.CLICKHOUSE_URL, - port: 8123, - basicAuth: { - username: process.env.CLICKHOUSE_USER, - password: process.env.CLICKHOUSE_PASSWORD - } +const clickhouse = createClient({ + url: process.env.CLICKHOUSE_URL ?? 'http://localhost:8123', + username: process.env.CLICKHOUSE_USER, + password: process.env.CLICKHOUSE_PASSWORD }) // PASS: バッチ挿入(効率的) async function bulkInsertTrades(trades: Trade[]) { - const values = trades.map(trade => `( - '${trade.id}', - '${trade.market_id}', - '${trade.user_id}', - ${trade.amount}, - '${trade.timestamp.toISOString()}' - )`).join(',') - - await clickhouse.query(` - INSERT INTO trades (id, market_id, user_id, amount, timestamp) - VALUES ${values} - `).toPromise() + await clickhouse.insert({ + table: 'trades', + values: trades.map(trade => ({ + id: trade.id, + market_id: trade.market_id, + user_id: trade.user_id, + amount: trade.amount, + timestamp: trade.timestamp.toISOString() + })), + format: 'JSONEachRow' + }) } // FAIL: 個別挿入(低速) async function insertTrade(trade: Trade) { // ループ内でこれをしないでください! - await clickhouse.query(` - INSERT INTO trades VALUES ('${trade.id}', ...) - `).toPromise() + await clickhouse.insert({ + table: 'trades', + values: [{ + id: trade.id, + market_id: trade.market_id, + user_id: trade.user_id, + amount: trade.amount, + timestamp: trade.timestamp.toISOString() + }], + format: 'JSONEachRow' + }) } ``` @@ -191,17 +195,14 @@ async function insertTrade(trade: Trade) { ```typescript // 継続的なデータ取り込み用 -import { createWriteStream } from 'fs' -import { pipeline } from 'stream/promises' +import { Readable } from 'node:stream' -async function streamInserts() { - const stream = clickhouse.insert('trades').stream() - - for await (const batch of dataSource) { - stream.write(batch) - } - - await stream.end() +async function streamInserts(dataSource: AsyncIterable>) { + await clickhouse.insert({ + table: 'trades', + values: Readable.from(dataSource, { objectMode: true }), + format: 'JSONEachRow' + }) } ``` @@ -386,14 +387,18 @@ pgClient.query('LISTEN market_updates') pgClient.on('notification', async (msg) => { const update = JSON.parse(msg.payload) - await clickhouse.insert('market_updates', [ - { - market_id: update.id, - event_type: update.operation, // INSERT, UPDATE, DELETE - timestamp: new Date(), - data: JSON.stringify(update.new_data) - } - ]) + await clickhouse.insert({ + table: 'market_updates', + values: [ + { + market_id: update.id, + event_type: update.operation, // INSERT, UPDATE, DELETE + timestamp: new Date(), + data: JSON.stringify(update.new_data) + } + ], + format: 'JSONEachRow' + }) }) ``` diff --git a/docs/ko-KR/skills/clickhouse-io/SKILL.md b/docs/ko-KR/skills/clickhouse-io/SKILL.md index 069604028..5d6b00805 100644 --- a/docs/ko-KR/skills/clickhouse-io/SKILL.md +++ b/docs/ko-KR/skills/clickhouse-io/SKILL.md @@ -161,36 +161,43 @@ ORDER BY market_id, date; ### 배치 삽입 (권장) ```typescript -import { ClickHouse } from 'clickhouse' +import { createClient } from '@clickhouse/client' -const clickhouse = new ClickHouse({ - url: process.env.CLICKHOUSE_URL, - port: 8123, - basicAuth: { - username: process.env.CLICKHOUSE_USER, - password: process.env.CLICKHOUSE_PASSWORD - } +const clickhouse = createClient({ + url: process.env.CLICKHOUSE_URL ?? 'http://localhost:8123', + username: process.env.CLICKHOUSE_USER, + password: process.env.CLICKHOUSE_PASSWORD }) // PASS: 배치 삽입 (효율적) async function bulkInsertTrades(trades: Trade[]) { - const rows = trades.map(trade => ({ - id: trade.id, - market_id: trade.market_id, - user_id: trade.user_id, - amount: trade.amount, - timestamp: trade.timestamp.toISOString() - })) - - await clickhouse.insert('trades', rows) + await clickhouse.insert({ + table: 'trades', + values: trades.map(trade => ({ + id: trade.id, + market_id: trade.market_id, + user_id: trade.user_id, + amount: trade.amount, + timestamp: trade.timestamp.toISOString() + })), + format: 'JSONEachRow' + }) } // FAIL: 개별 삽입 (느림) async function insertTrade(trade: Trade) { // 루프 안에서 이렇게 하지 마세요! - await clickhouse.query(` - INSERT INTO trades VALUES ('${trade.id}', ...) - `).toPromise() + await clickhouse.insert({ + table: 'trades', + values: [{ + id: trade.id, + market_id: trade.market_id, + user_id: trade.user_id, + amount: trade.amount, + timestamp: trade.timestamp.toISOString() + }], + format: 'JSONEachRow' + }) } ``` @@ -198,17 +205,14 @@ async function insertTrade(trade: Trade) { ```typescript // 연속적인 데이터 수집용 -import { createWriteStream } from 'fs' -import { pipeline } from 'stream/promises' +import { Readable } from 'node:stream' -async function streamInserts() { - const stream = clickhouse.insert('trades').stream() - - for await (const batch of dataSource) { - stream.write(batch) - } - - await stream.end() +async function streamInserts(dataSource: AsyncIterable>) { + await clickhouse.insert({ + table: 'trades', + values: Readable.from(dataSource, { objectMode: true }), + format: 'JSONEachRow' + }) } ``` @@ -404,14 +408,18 @@ pgClient.query('LISTEN market_updates') pgClient.on('notification', async (msg) => { const update = JSON.parse(msg.payload) - await clickhouse.insert('market_updates', [ - { - market_id: update.id, - event_type: update.operation, // INSERT, UPDATE, DELETE - timestamp: new Date(), - data: JSON.stringify(update.new_data) - } - ]) + await clickhouse.insert({ + table: 'market_updates', + values: [ + { + market_id: update.id, + event_type: update.operation, // INSERT, UPDATE, DELETE + timestamp: new Date(), + data: JSON.stringify(update.new_data) + } + ], + format: 'JSONEachRow' + }) }) ``` diff --git a/docs/zh-CN/skills/clickhouse-io/SKILL.md b/docs/zh-CN/skills/clickhouse-io/SKILL.md index fd929aabf..52c01da43 100644 --- a/docs/zh-CN/skills/clickhouse-io/SKILL.md +++ b/docs/zh-CN/skills/clickhouse-io/SKILL.md @@ -162,39 +162,43 @@ ORDER BY market_id, date; ### 批量插入 (推荐) ```typescript -import { ClickHouse } from 'clickhouse' +import { createClient } from '@clickhouse/client' -const clickhouse = new ClickHouse({ - url: process.env.CLICKHOUSE_URL, - port: 8123, - basicAuth: { - username: process.env.CLICKHOUSE_USER, - password: process.env.CLICKHOUSE_PASSWORD - } +const clickhouse = createClient({ + url: process.env.CLICKHOUSE_URL ?? 'http://localhost:8123', + username: process.env.CLICKHOUSE_USER, + password: process.env.CLICKHOUSE_PASSWORD }) // PASS: Batch insert (efficient) async function bulkInsertTrades(trades: Trade[]) { - const values = trades.map(trade => `( - '${trade.id}', - '${trade.market_id}', - '${trade.user_id}', - ${trade.amount}, - '${trade.timestamp.toISOString()}' - )`).join(',') - - await clickhouse.query(` - INSERT INTO trades (id, market_id, user_id, amount, timestamp) - VALUES ${values} - `).toPromise() + await clickhouse.insert({ + table: 'trades', + values: trades.map(trade => ({ + id: trade.id, + market_id: trade.market_id, + user_id: trade.user_id, + amount: trade.amount, + timestamp: trade.timestamp.toISOString() + })), + format: 'JSONEachRow' + }) } // FAIL: Individual inserts (slow) async function insertTrade(trade: Trade) { // Don't do this in a loop! - await clickhouse.query(` - INSERT INTO trades VALUES ('${trade.id}', ...) - `).toPromise() + await clickhouse.insert({ + table: 'trades', + values: [{ + id: trade.id, + market_id: trade.market_id, + user_id: trade.user_id, + amount: trade.amount, + timestamp: trade.timestamp.toISOString() + }], + format: 'JSONEachRow' + }) } ``` @@ -202,17 +206,14 @@ async function insertTrade(trade: Trade) { ```typescript // For continuous data ingestion -import { createWriteStream } from 'fs' -import { pipeline } from 'stream/promises' +import { Readable } from 'node:stream' -async function streamInserts() { - const stream = clickhouse.insert('trades').stream() - - for await (const batch of dataSource) { - stream.write(batch) - } - - await stream.end() +async function streamInserts(dataSource: AsyncIterable>) { + await clickhouse.insert({ + table: 'trades', + values: Readable.from(dataSource, { objectMode: true }), + format: 'JSONEachRow' + }) } ``` @@ -397,14 +398,18 @@ pgClient.query('LISTEN market_updates') pgClient.on('notification', async (msg) => { const update = JSON.parse(msg.payload) - await clickhouse.insert('market_updates', [ - { - market_id: update.id, - event_type: update.operation, // INSERT, UPDATE, DELETE - timestamp: new Date(), - data: JSON.stringify(update.new_data) - } - ]) + await clickhouse.insert({ + table: 'market_updates', + values: [ + { + market_id: update.id, + event_type: update.operation, // INSERT, UPDATE, DELETE + timestamp: new Date(), + data: JSON.stringify(update.new_data) + } + ], + format: 'JSONEachRow' + }) }) ``` diff --git a/docs/zh-TW/skills/clickhouse-io/SKILL.md b/docs/zh-TW/skills/clickhouse-io/SKILL.md index aaa95d4fc..913c72e84 100644 --- a/docs/zh-TW/skills/clickhouse-io/SKILL.md +++ b/docs/zh-TW/skills/clickhouse-io/SKILL.md @@ -151,39 +151,43 @@ ORDER BY market_id, date; ### 批量插入(推薦) ```typescript -import { ClickHouse } from 'clickhouse' +import { createClient } from '@clickhouse/client' -const clickhouse = new ClickHouse({ - url: process.env.CLICKHOUSE_URL, - port: 8123, - basicAuth: { - username: process.env.CLICKHOUSE_USER, - password: process.env.CLICKHOUSE_PASSWORD - } +const clickhouse = createClient({ + url: process.env.CLICKHOUSE_URL ?? 'http://localhost:8123', + username: process.env.CLICKHOUSE_USER, + password: process.env.CLICKHOUSE_PASSWORD }) // PASS: 批量插入(高效) async function bulkInsertTrades(trades: Trade[]) { - const values = trades.map(trade => `( - '${trade.id}', - '${trade.market_id}', - '${trade.user_id}', - ${trade.amount}, - '${trade.timestamp.toISOString()}' - )`).join(',') - - await clickhouse.query(` - INSERT INTO trades (id, market_id, user_id, amount, timestamp) - VALUES ${values} - `).toPromise() + await clickhouse.insert({ + table: 'trades', + values: trades.map(trade => ({ + id: trade.id, + market_id: trade.market_id, + user_id: trade.user_id, + amount: trade.amount, + timestamp: trade.timestamp.toISOString() + })), + format: 'JSONEachRow' + }) } // FAIL: 個別插入(慢) async function insertTrade(trade: Trade) { // 不要在迴圈中這樣做! - await clickhouse.query(` - INSERT INTO trades VALUES ('${trade.id}', ...) - `).toPromise() + await clickhouse.insert({ + table: 'trades', + values: [{ + id: trade.id, + market_id: trade.market_id, + user_id: trade.user_id, + amount: trade.amount, + timestamp: trade.timestamp.toISOString() + }], + format: 'JSONEachRow' + }) } ``` @@ -191,17 +195,14 @@ async function insertTrade(trade: Trade) { ```typescript // 用於持續資料攝取 -import { createWriteStream } from 'fs' -import { pipeline } from 'stream/promises' +import { Readable } from 'node:stream' -async function streamInserts() { - const stream = clickhouse.insert('trades').stream() - - for await (const batch of dataSource) { - stream.write(batch) - } - - await stream.end() +async function streamInserts(dataSource: AsyncIterable>) { + await clickhouse.insert({ + table: 'trades', + values: Readable.from(dataSource, { objectMode: true }), + format: 'JSONEachRow' + }) } ``` @@ -386,14 +387,18 @@ pgClient.query('LISTEN market_updates') pgClient.on('notification', async (msg) => { const update = JSON.parse(msg.payload) - await clickhouse.insert('market_updates', [ - { - market_id: update.id, - event_type: update.operation, // INSERT, UPDATE, DELETE - timestamp: new Date(), - data: JSON.stringify(update.new_data) - } - ]) + await clickhouse.insert({ + table: 'market_updates', + values: [ + { + market_id: update.id, + event_type: update.operation, // INSERT, UPDATE, DELETE + timestamp: new Date(), + data: JSON.stringify(update.new_data) + } + ], + format: 'JSONEachRow' + }) }) ``` diff --git a/skills/clickhouse-io/SKILL.md b/skills/clickhouse-io/SKILL.md index 9753e59f3..a0bc18f79 100644 --- a/skills/clickhouse-io/SKILL.md +++ b/skills/clickhouse-io/SKILL.md @@ -162,39 +162,43 @@ ORDER BY market_id, date; ### Bulk Insert (Recommended) ```typescript -import { ClickHouse } from 'clickhouse' +import { createClient } from '@clickhouse/client' -const clickhouse = new ClickHouse({ - url: process.env.CLICKHOUSE_URL, - port: 8123, - basicAuth: { - username: process.env.CLICKHOUSE_USER, - password: process.env.CLICKHOUSE_PASSWORD - } +const clickhouse = createClient({ + url: process.env.CLICKHOUSE_URL ?? 'http://localhost:8123', + username: process.env.CLICKHOUSE_USER, + password: process.env.CLICKHOUSE_PASSWORD }) // PASS: Batch insert (efficient) async function bulkInsertTrades(trades: Trade[]) { - const values = trades.map(trade => `( - '${trade.id}', - '${trade.market_id}', - '${trade.user_id}', - ${trade.amount}, - '${trade.timestamp.toISOString()}' - )`).join(',') - - await clickhouse.query(` - INSERT INTO trades (id, market_id, user_id, amount, timestamp) - VALUES ${values} - `).toPromise() + await clickhouse.insert({ + table: 'trades', + values: trades.map(trade => ({ + id: trade.id, + market_id: trade.market_id, + user_id: trade.user_id, + amount: trade.amount, + timestamp: trade.timestamp.toISOString() + })), + format: 'JSONEachRow' + }) } // FAIL: Individual inserts (slow) async function insertTrade(trade: Trade) { // Don't do this in a loop! - await clickhouse.query(` - INSERT INTO trades VALUES ('${trade.id}', ...) - `).toPromise() + await clickhouse.insert({ + table: 'trades', + values: [{ + id: trade.id, + market_id: trade.market_id, + user_id: trade.user_id, + amount: trade.amount, + timestamp: trade.timestamp.toISOString() + }], + format: 'JSONEachRow' + }) } ``` @@ -202,17 +206,14 @@ async function insertTrade(trade: Trade) { ```typescript // For continuous data ingestion -import { createWriteStream } from 'fs' -import { pipeline } from 'stream/promises' +import { Readable } from 'node:stream' -async function streamInserts() { - const stream = clickhouse.insert('trades').stream() - - for await (const batch of dataSource) { - stream.write(batch) - } - - await stream.end() +async function streamInserts(dataSource: AsyncIterable>) { + await clickhouse.insert({ + table: 'trades', + values: Readable.from(dataSource, { objectMode: true }), + format: 'JSONEachRow' + }) } ``` @@ -397,14 +398,18 @@ pgClient.query('LISTEN market_updates') pgClient.on('notification', async (msg) => { const update = JSON.parse(msg.payload) - await clickhouse.insert('market_updates', [ - { - market_id: update.id, - event_type: update.operation, // INSERT, UPDATE, DELETE - timestamp: new Date(), - data: JSON.stringify(update.new_data) - } - ]) + await clickhouse.insert({ + table: 'market_updates', + values: [ + { + market_id: update.id, + event_type: update.operation, // INSERT, UPDATE, DELETE + timestamp: new Date(), + data: JSON.stringify(update.new_data) + } + ], + format: 'JSONEachRow' + }) }) ```