Skip to content

Commit 09e9585

Browse files
refactor: remove DeltaLake dependency, replace with ClickHouse pipeline
DeltaLake was never fully wired up (factory files never created). Now that events flow through Pipeline → R2 → S3Queue → ClickHouse, remove the dead DeltaLake code paths: - Remove writeEventsDelta/writeEventsDualWrite from event-writer.ts - Remove USE_DELTALAKE/USE_DELTALAKE_CDC/DELTALAKE_DUAL_WRITE env vars - Remove DeltaLake compaction tasks from scheduled handler - Simplify EventWriterDO flush to legacy Parquet only - Add missing isCollectionChangeEvent export to cdc-processor.ts - Add headlessly-tail to tail_consumers Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
1 parent 7f83c15 commit 09e9585

6 files changed

Lines changed: 16 additions & 413 deletions

File tree

core/src/cdc-processor.ts

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,13 @@ import {
2323
export type { CDCEvent } from './cdc-delta.js'
2424
export { isCdcEvent } from './types.js'
2525

26+
/**
27+
* Type guard for collection change events (CDC: insert/update/delete).
28+
*/
29+
export function isCollectionChangeEvent(event: { type?: string }): boolean {
30+
return typeof event.type === 'string' && event.type.startsWith('collection.')
31+
}
32+
2633
// ============================================================================
2734
// Types
2835
// ============================================================================

src/env.ts

Lines changed: 0 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -99,14 +99,4 @@ export interface Env extends WebhookEnv {
9999
/** Max request body size in bytes for query endpoint (default: 65536 = 64KB) */
100100
MAX_QUERY_BODY_SIZE?: string
101101

102-
// ============================================================================
103-
// DeltaLake Integration Configuration
104-
// ============================================================================
105-
106-
/** Enable DeltaTable for event storage (default: false) */
107-
USE_DELTALAKE?: string
108-
/** Enable DeltaTable for CDC storage (default: false) */
109-
USE_DELTALAKE_CDC?: string
110-
/** Enable dual-write mode during migration - writes to both legacy and DeltaTable (default: false) */
111-
DELTALAKE_DUAL_WRITE?: string
112102
}

src/event-writer-do.ts

Lines changed: 6 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,7 @@
1515
*/
1616

1717
import { DurableObject } from 'cloudflare:workers'
18-
import { writeEvents, writeEventsDelta, writeEventsDualWrite, ulid, type EventRecord, type WriteResult, type DeltaWriteResult } from './event-writer'
18+
import { writeEvents, ulid, type EventRecord, type WriteResult } from './event-writer'
1919
import { recordWriterDOMetric, recordR2WriteMetric, MetricTimer } from './metrics'
2020
import type { ShardCoordinatorDO } from './shard-coordinator-do'
2121
import { logger, sanitize, logError, type Logger } from './logger'
@@ -41,7 +41,7 @@ const DEFAULT_DEDUP_TTL_MS = 24 * 60 * 60 * 1000 // 24 hours default TTL for ded
4141

4242
import type { Env as FullEnv } from './env'
4343

44-
export type Env = Pick<FullEnv, 'EVENTS_BUCKET' | 'EVENT_WRITER' | 'ANALYTICS' | 'SHARD_COORDINATOR' | 'USE_DELTALAKE' | 'DELTALAKE_DUAL_WRITE'>
44+
export type Env = Pick<FullEnv, 'EVENTS_BUCKET' | 'EVENT_WRITER' | 'ANALYTICS' | 'SHARD_COORDINATOR'>
4545

4646
// ============================================================================
4747
// Configuration
@@ -547,38 +547,15 @@ export class EventWriterDO extends DurableObject<Env> {
547547

548548
const results: WriteResult[] = []
549549
let totalBytes = 0
550-
const useDeltalake = this.env.USE_DELTALAKE === 'true'
551-
const useDualWrite = this.env.DELTALAKE_DUAL_WRITE === 'true'
552550
try {
553551
// Write each source group to its own prefix
554552
for (const [source, sourceEvents] of eventsBySource) {
555553
const writeTimer = new MetricTimer()
556-
let result: WriteResult
557-
let bytes: number
558-
559-
if (useDualWrite) {
560-
// Dual-write mode: write to both legacy and DeltaTable during migration
561-
const dualResult = await writeEventsDualWrite(this.env.EVENTS_BUCKET, source, sourceEvents, this.shardId)
562-
result = dualResult.legacy
563-
bytes = dualResult.legacy.bytes + dualResult.delta.bytes
564-
} else if (useDeltalake) {
565-
// DeltaTable-only mode
566-
const deltaResult = await writeEventsDelta(this.env.EVENTS_BUCKET, sourceEvents, this.shardId)
567-
// Convert delta result to WriteResult format for compatibility
568-
result = {
569-
key: deltaResult.path,
570-
bytes: deltaResult.bytes,
571-
events: deltaResult.events,
572-
cpuMs: deltaResult.cpuMs,
573-
}
574-
bytes = deltaResult.bytes
575-
} else {
576-
// Legacy mode (default)
577-
result = await writeEvents(this.env.EVENTS_BUCKET, source, sourceEvents)
578-
bytes = result.bytes
579-
}
580554

581-
logger.child({ component: 'EventWriterDO', shard: this.shardId }).info('Flushed', { key: sanitize.id(result.key, 64), mode: useDualWrite ? 'dual' : useDeltalake ? 'deltalake' : 'legacy' })
555+
const result = await writeEvents(this.env.EVENTS_BUCKET, source, sourceEvents)
556+
const bytes = result.bytes
557+
558+
logger.child({ component: 'EventWriterDO', shard: this.shardId }).info('Flushed', { key: sanitize.id(result.key, 64) })
582559
results.push(result)
583560
totalBytes += bytes
584561

src/event-writer.ts

Lines changed: 1 addition & 153 deletions
Original file line numberDiff line numberDiff line change
@@ -11,11 +11,8 @@
1111
*/
1212

1313
import { parquetWriteBuffer } from '@dotdo/hyparquet-writer'
14-
import { DeltaTable, type StorageBackend } from '@dotdo/deltalake'
1514
import { ulid } from '../core/src/ulid'
16-
import { sanitizeR2Path, buildSafeR2Path, sanitizePathSegment, InvalidR2PathError } from './utils'
17-
import { createR2Storage } from '../core/src/deltalake-storage'
18-
import { createEventsTable, getEventsTablePath, type EventDeltaRecord } from '../core/src/deltalake-factory'
15+
import { sanitizeR2Path } from './utils'
1916

2017
export { ulid }
2118

@@ -159,155 +156,6 @@ export async function writeEvents(
159156
return { key, events: events.length, bytes: buffer.byteLength, cpuMs }
160157
}
161158

162-
// ============================================================================
163-
// DeltaLake Event Writer
164-
// ============================================================================
165-
166-
/**
167-
* Result from a DeltaLake write operation.
168-
*/
169-
export interface DeltaWriteResult {
170-
version: number
171-
events: number
172-
bytes: number
173-
path: string
174-
cpuMs: number
175-
}
176-
177-
/**
178-
* Convert EventRecord to DeltaLake format (EventDeltaRecord).
179-
*/
180-
function toEventDeltaRecord(event: EventRecord): EventDeltaRecord {
181-
return {
182-
ts: event.ts,
183-
type: event.type,
184-
source: event.source ?? null,
185-
provider: event.provider ?? null,
186-
eventType: event.eventType ?? null,
187-
verified: event.verified ?? null,
188-
scriptName: event.scriptName ?? null,
189-
outcome: event.outcome ?? null,
190-
method: event.method ?? null,
191-
url: event.url ?? null,
192-
statusCode: event.statusCode ?? null,
193-
durationMs: event.durationMs ?? null,
194-
payload: event.payload ? JSON.stringify(event.payload) : null,
195-
}
196-
}
197-
198-
/**
199-
* DeltaLake table cache - one table per shard.
200-
* Maps shardId to DeltaTable instance.
201-
*/
202-
const deltaTableCache = new Map<string, DeltaTable<EventDeltaRecord>>()
203-
204-
/**
205-
* Get or create a DeltaTable for a shard.
206-
*/
207-
function getDeltaTable(
208-
storage: StorageBackend,
209-
shardId: number
210-
): DeltaTable<EventDeltaRecord> {
211-
const cacheKey = `events-shard-${shardId}`
212-
let table = deltaTableCache.get(cacheKey)
213-
214-
if (!table) {
215-
table = createEventsTable(storage, { shardId })
216-
deltaTableCache.set(cacheKey, table)
217-
}
218-
219-
return table
220-
}
221-
222-
/**
223-
* Write events to R2 using DeltaLake for ACID guarantees.
224-
*
225-
* Events are written to a partitioned DeltaTable with transaction log support.
226-
* Path structure: events/shard={shardId}/_delta_log/...
227-
*
228-
* @param bucket - R2 bucket for storage
229-
* @param events - Events to write
230-
* @param shardId - Shard ID for partitioning (default: 0)
231-
* @returns Result with commit version and event count
232-
*
233-
* @example
234-
* ```typescript
235-
* const result = await writeEventsDelta(env.EVENTS_BUCKET, events, 2)
236-
* console.log(`Committed version ${result.version} with ${result.events} events`)
237-
* ```
238-
*/
239-
export async function writeEventsDelta(
240-
bucket: R2Bucket,
241-
events: EventRecord[],
242-
shardId: number = 0
243-
): Promise<DeltaWriteResult> {
244-
if (events.length === 0) {
245-
throw new Error('No events to write')
246-
}
247-
248-
const startCpu = performance.now()
249-
250-
// Create storage backend from R2 bucket
251-
const storage = createR2Storage(bucket)
252-
253-
// Get or create the DeltaTable for this shard
254-
const table = getDeltaTable(storage, shardId)
255-
256-
// Convert events to DeltaLake format
257-
const deltaRecords = events.map(toEventDeltaRecord)
258-
259-
// Estimate bytes from serialized records (for metrics)
260-
const estimatedBytes = deltaRecords.reduce((sum, record) => {
261-
return sum + JSON.stringify(record).length
262-
}, 0)
263-
264-
// Write with ACID guarantees
265-
const commit = await table.write(deltaRecords)
266-
267-
const cpuMs = performance.now() - startCpu
268-
const tablePath = getEventsTablePath(shardId)
269-
270-
console.log(`[WRITE-DELTA] shard=${shardId} version=${commit.version} events=${events.length} cpuMs=${cpuMs.toFixed(2)}`)
271-
272-
return {
273-
version: commit.version,
274-
events: events.length,
275-
bytes: estimatedBytes,
276-
path: `${tablePath}/_delta_log/${String(commit.version).padStart(20, '0')}.json`,
277-
cpuMs,
278-
}
279-
}
280-
281-
/**
282-
* Write events using both legacy Parquet and DeltaLake (dual-write mode).
283-
*
284-
* This function writes to both storage formats during migration, ensuring
285-
* data consistency and enabling gradual rollout of DeltaLake.
286-
*
287-
* @param bucket - R2 bucket for storage
288-
* @param prefix - Legacy path prefix
289-
* @param events - Events to write
290-
* @param shardId - Shard ID for DeltaLake partitioning
291-
* @returns Both legacy and DeltaLake results
292-
*/
293-
export async function writeEventsDualWrite(
294-
bucket: R2Bucket,
295-
prefix: string,
296-
events: EventRecord[],
297-
shardId: number = 0
298-
): Promise<{ legacy: WriteResult; delta: DeltaWriteResult }> {
299-
// Write to both in parallel for performance
300-
const [legacyResult, deltaResult] = await Promise.all([
301-
writeEvents(bucket, prefix, events),
302-
writeEventsDelta(bucket, events, shardId),
303-
])
304-
305-
return {
306-
legacy: legacyResult,
307-
delta: deltaResult,
308-
}
309-
}
310-
311159
// ============================================================================
312160
// Event Buffer - For batching events before writing
313161
// ============================================================================

0 commit comments

Comments
 (0)