Skip to content
This repository was archived by the owner on Nov 4, 2021. It is now read-only.
Merged
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
2 changes: 1 addition & 1 deletion package.json
Original file line number Diff line number Diff line change
Expand Up @@ -37,10 +37,10 @@
"license": "MIT",
"dependencies": {
"@google-cloud/bigquery": "^5.5.0",
"@posthog/clickhouse": "^1.7.0",
"@sentry/node": "^5.29.0",
"@sentry/tracing": "^5.29.0",
"adm-zip": "^0.4.16",
"clickhouse": "^2.2.1",
"fastify": "^3.8.0",
"hot-shots": "^8.2.1",
"ioredis": "^4.19.2",
Expand Down
41 changes: 23 additions & 18 deletions src/db.ts
Original file line number Diff line number Diff line change
@@ -1,9 +1,8 @@
import { Properties } from '@posthog/plugin-scaffold'
import { ClickHouse } from 'clickhouse'
import ClickHouse from '@posthog/clickhouse'
import { Producer } from 'kafkajs'
import { DateTime } from 'luxon'
import { Pool, QueryConfig, QueryResult, QueryResultRow } from 'pg'
import { string } from 'yargs'
import { KAFKA_PERSON, KAFKA_PERSON_UNIQUE_ID } from './ingestion/topics'
import { chainToElements, hashElements, unparsePersonPartial } from './ingestion/utils'
import {
Expand Down Expand Up @@ -45,14 +44,17 @@ export class DB {
queryTextOrConfig: string | QueryConfig<I>,
values?: I
): Promise<QueryResult<R>> {
return this.postgres.query(queryTextOrConfig, values)
return await this.postgres.query(queryTextOrConfig, values)
}

public async clickhouseQuery(query: string, reqParams?: Record<string, any>): Promise<Record<string, any>> {
public async clickhouseQuery(
query: string,
options?: ClickHouse.QueryOptions
): Promise<ClickHouse.QueryResult<Record<string, any>>> {
if (!this.clickhouse) {
throw new Error('ClickHouse connection has not been provided to this DB instance!')
}
return this.clickhouse.query(query, reqParams).toPromise()
return await this.clickhouse.querying(query, options)
}

// Person
Expand All @@ -61,7 +63,7 @@ export class DB {
public async fetchPersons(database: Database.ClickHouse): Promise<ClickHousePerson[]>
public async fetchPersons(database: Database = Database.Postgres): Promise<Person[] | ClickHousePerson[]> {
if (database === Database.ClickHouse) {
return (await this.clickhouseQuery('SELECT * FROM person')) as ClickHousePerson[]
return (await this.clickhouseQuery('SELECT * FROM person')).data as ClickHousePerson[]
Comment thread
mariusandra marked this conversation as resolved.
} else if (database === Database.Postgres) {
return ((await this.postgresQuery('SELECT * FROM posthog_person')).rows as RawPerson[]).map(
(rawPerson: RawPerson) =>
Expand Down Expand Up @@ -186,11 +188,13 @@ export class DB {
database: Database = Database.Postgres
): Promise<PersonDistinctId[] | ClickHousePersonDistinctId[]> {
if (database === Database.ClickHouse) {
return (await this.clickhouseQuery(
`SELECT * FROM person_distinct_id WHERE person_id='${escapeClickHouseString(
person.uuid
)}' and team_id='${person.team_id}' ORDER BY id`
)) as ClickHousePersonDistinctId[]
return (
await this.clickhouseQuery(
`SELECT * FROM person_distinct_id WHERE person_id='${escapeClickHouseString(
person.uuid
)}' and team_id='${person.team_id}' ORDER BY id`
)
).data as ClickHousePersonDistinctId[]
} else if (database === Database.Postgres) {
const result = await this.postgresQuery(
'SELECT * FROM posthog_persondistinctid WHERE person_id=$1 and team_id=$2 ORDER BY id',
Expand Down Expand Up @@ -255,7 +259,7 @@ export class DB {

public async fetchEvents(): Promise<Event[] | ClickHouseEvent[]> {
if (this.kafkaProducer) {
const events = (await this.clickhouseQuery(`SELECT * FROM events`)) as ClickHouseEvent[]
const events = (await this.clickhouseQuery(`SELECT * FROM events`)).data as ClickHouseEvent[]
return (
events?.map(
(event) =>
Expand All @@ -278,9 +282,8 @@ export class DB {

public async fetchSessionRecordingEvents(): Promise<PostgresSessionRecordingEvent[] | SessionRecordingEvent[]> {
if (this.kafkaProducer) {
const events = ((await this.clickhouseQuery(
`SELECT * FROM session_recording_events`
)) as SessionRecordingEvent[]).map((event) => {
const events = ((await this.clickhouseQuery(`SELECT * FROM session_recording_events`))
.data as SessionRecordingEvent[]).map((event) => {
return {
...event,
snapshot_data: event.snapshot_data ? JSON.parse(event.snapshot_data) : null,
Expand All @@ -297,9 +300,11 @@ export class DB {

public async fetchElements(event?: Event): Promise<Element[]> {
if (this.kafkaProducer) {
const events = (await this.clickhouseQuery(
`SELECT elements_chain FROM events WHERE uuid='${escapeClickHouseString((event as any).uuid)}'`
)) as ClickHouseEvent[]
const events = (
await this.clickhouseQuery(
`SELECT elements_chain FROM events WHERE uuid='${escapeClickHouseString((event as any).uuid)}'`
)
).data as ClickHouseEvent[]
const chain = events?.[0]?.elements_chain
return chainToElements(chain)
} else {
Expand Down
2 changes: 1 addition & 1 deletion src/ingestion/process-event.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ import { Event as EventProto, IEvent } from '../idl/protos'
import { Producer } from 'kafkajs'
import { KAFKA_EVENTS, KAFKA_SESSION_RECORDING_EVENTS } from './topics'
import { elementsToString, sanitizeEventName } from './utils'
import { ClickHouse } from 'clickhouse'
import ClickHouse from '@posthog/clickhouse'
import { DB } from '../db'
import { status } from '../status'
import * as Sentry from '@sentry/node'
Expand Down
19 changes: 11 additions & 8 deletions src/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ import { Kafka, logLevel, Producer } from 'kafkajs'
import { FastifyInstance } from 'fastify'
import { PluginsServer, PluginsServerConfig, Queue } from './types'
import { startQueue } from './worker/queue'
import { ClickHouse } from 'clickhouse'
import ClickHouse from '@posthog/clickhouse'
import { startFastifyInstance, stopFastifyInstance } from './web/server'
import { version } from '../package.json'
import { PluginEvent } from '@posthog/plugin-scaffold'
Expand All @@ -20,6 +20,7 @@ import { startSchedule } from './services/schedule'
import { ConnectionOptions } from 'tls'
import { DB } from './db'
import { DateTime } from 'luxon'
import * as fs from 'fs'
import { KAFKA_EVENTS_PLUGIN_INGESTION, KAFKA_EVENTS_WAL } from './ingestion/topics'

export async function createServer(
Expand Down Expand Up @@ -71,17 +72,19 @@ export async function createServer(
throw new Error('You must set KAFKA_HOSTS to process events from Kafka!')
}
clickhouse = new ClickHouse({
url: `http${serverConfig.CLICKHOUSE_SECURE ? 's' : ''}://$${serverConfig.CLICKHOUSE_HOST}`,
host: serverConfig.CLICKHOUSE_HOST,
port: serverConfig.CLICKHOUSE_SECURE ? 8443 : 8123,
basicAuth: {
username: serverConfig.CLICKHOUSE_USER,
password: serverConfig.CLICKHOUSE_PASSWORD,
},
config: {
protocol: serverConfig.CLICKHOUSE_SECURE ? 'https:' : 'http:',
user: serverConfig.CLICKHOUSE_USER,
password: serverConfig.CLICKHOUSE_PASSWORD || undefined,
dataObjects: true,
queryOptions: {
database: serverConfig.CLICKHOUSE_DATABASE,
output_format_json_quote_64bit_integers: false,
},
ca: serverConfig.CLICKHOUSE_CA ? fs.readFileSync(serverConfig.CLICKHOUSE_CA).toString() : undefined,
})
await clickhouse.query('SELECT 1') // test that the connection works
await clickhouse.querying('SELECT 1') // test that the connection works

if (!serverConfig.KAFKA_CONSUMPTION_TOPIC) {
// When ingesting events, listen to the "INGESTION_HANDOFF" topic, otherwise listen to the "WAL" and discard
Expand Down
2 changes: 1 addition & 1 deletion src/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ import { VM } from 'vm2'
import { DateTime } from 'luxon'
import { StatsD } from 'hot-shots'
import { EventsProcessor } from 'ingestion/process-event'
import { ClickHouse } from 'clickhouse'
import ClickHouse from '@posthog/clickhouse'
import { DB } from './db'

export enum LogLevel {
Expand Down
24 changes: 13 additions & 11 deletions tests/helpers/clickhouse.ts
Original file line number Diff line number Diff line change
@@ -1,22 +1,24 @@
import { defaultConfig } from '../../src/config'
import { ClickHouse } from 'clickhouse'
import ClickHouse from '@posthog/clickhouse'
import { PluginsServerConfig } from '../../src/types'

export async function resetTestDatabaseClickhouse(extraServerConfig: Partial<PluginsServerConfig>): Promise<void> {
const config = { ...defaultConfig, ...extraServerConfig }
const clickhouse = new ClickHouse({
url: `http://$${config.CLICKHOUSE_HOST}`,
host: config.CLICKHOUSE_HOST,
port: 8123,
config: {
dataObjects: true,
queryOptions: {
database: config.CLICKHOUSE_DATABASE,
output_format_json_quote_64bit_integers: false,
},
})
await clickhouse.query('TRUNCATE events').toPromise()
await clickhouse.query('TRUNCATE events_mv').toPromise()
await clickhouse.query('TRUNCATE person').toPromise()
await clickhouse.query('TRUNCATE person_distinct_id').toPromise()
await clickhouse.query('TRUNCATE person_mv').toPromise()
await clickhouse.query('TRUNCATE person_static_cohort').toPromise()
await clickhouse.query('TRUNCATE session_recording_events').toPromise()
await clickhouse.query('TRUNCATE session_recording_events_mv').toPromise()
await clickhouse.querying('TRUNCATE events')
await clickhouse.querying('TRUNCATE events_mv')
await clickhouse.querying('TRUNCATE person')
await clickhouse.querying('TRUNCATE person_distinct_id')
await clickhouse.querying('TRUNCATE person_mv')
await clickhouse.querying('TRUNCATE person_static_cohort')
await clickhouse.querying('TRUNCATE session_recording_events')
await clickhouse.querying('TRUNCATE session_recording_events_mv')
}
Loading