mirror of
https://github.com/affaan-m/ECC.git
synced 2026-08-17 21:15:40 +02:00
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.
This commit is contained in:
@@ -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<Record<string, unknown>>) {
|
||||
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'
|
||||
})
|
||||
})
|
||||
```
|
||||
|
||||
|
||||
@@ -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<Record<string, unknown>>) {
|
||||
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'
|
||||
})
|
||||
})
|
||||
```
|
||||
|
||||
|
||||
@@ -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<Record<string, unknown>>) {
|
||||
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'
|
||||
})
|
||||
})
|
||||
```
|
||||
|
||||
|
||||
@@ -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<Record<string, unknown>>) {
|
||||
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'
|
||||
})
|
||||
})
|
||||
```
|
||||
|
||||
|
||||
@@ -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<Record<string, unknown>>) {
|
||||
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'
|
||||
})
|
||||
})
|
||||
```
|
||||
|
||||
|
||||
Reference in New Issue
Block a user