ingest.ts 5.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166
  1. import { Buffer } from "node:buffer"
  2. import { FirehoseClient, PutRecordBatchCommand } from "@aws-sdk/client-firehose"
  3. import { Effect, Layer } from "effect"
  4. import * as Context from "effect/Context"
  5. import { Resource } from "sst/resource"
  6. const MAX_FIREHOSE_BATCH_SIZE = 500
  7. const MAX_FIREHOSE_ATTEMPTS = 3
  8. const LAKE_TYPE = /^([A-Za-z0-9_]+)\.([A-Za-z0-9_]+)$/
  9. type IngestEvent = Record<string, unknown>
  10. type LakeRoute = { database: string; table: string }
  11. type FirehoseRecord = { Data: Uint8Array }
  12. export class IngestError extends Error {
  13. readonly _tag = "IngestError"
  14. readonly failed: number
  15. constructor(input: { message: string; failed: number; cause?: unknown }) {
  16. super(input.message, { cause: input.cause })
  17. this.name = "IngestError"
  18. this.failed = input.failed
  19. }
  20. }
  21. export declare namespace Ingest {
  22. export interface Service {
  23. readonly write: (events: unknown[]) => Effect.Effect<{ records: number }, IngestError>
  24. }
  25. }
  26. export class Ingest extends Context.Service<Ingest, Ingest.Service>()("@kirincode/stats/Ingest") {
  27. static readonly layer: Layer.Layer<Ingest> = Layer.effect(
  28. Ingest,
  29. Effect.sync(() => {
  30. const client = new FirehoseClient({})
  31. const write = Effect.fn("Ingest.write")(function* (events: unknown[]) {
  32. if (events.length === 0) return { records: 0 }
  33. const counts = countRoutedEvents(events)
  34. if (counts.unsupported > 0) {
  35. yield* Effect.logWarning(
  36. `lake ingest rejected ${JSON.stringify({ records: counts.records, unsupported: counts.unsupported })}`,
  37. )
  38. return yield* Effect.fail(
  39. new IngestError({
  40. message: "Unsupported lake event type",
  41. failed: counts.unsupported,
  42. }),
  43. )
  44. }
  45. if (counts.records === 0) return { records: 0 }
  46. let batch: FirehoseRecord[] = []
  47. let batches = 0
  48. let failed = 0
  49. for (const event of events) {
  50. if (!isRecord(event)) continue
  51. const route = routeEvent(event)
  52. if (!route) continue
  53. batch.push(toFirehoseRecord(event, route))
  54. if (batch.length < MAX_FIREHOSE_BATCH_SIZE) continue
  55. failed += yield* putRecords(client, Resource.LakeIngestConfig.streamName, batch)
  56. batches++
  57. batch = []
  58. }
  59. if (batch.length > 0) {
  60. failed += yield* putRecords(client, Resource.LakeIngestConfig.streamName, batch)
  61. batches++
  62. }
  63. if (failed > 0) {
  64. yield* Effect.logWarning(`lake ingest incomplete ${JSON.stringify({ records: counts.records, failed })}`)
  65. return yield* Effect.fail(new IngestError({ message: "Failed to ingest all lake records", failed }))
  66. }
  67. yield* Effect.logInfo(`lake ingest complete ${JSON.stringify({ records: counts.records, batches })}`)
  68. return { records: counts.records }
  69. })
  70. return Ingest.of({ write })
  71. }),
  72. )
  73. }
  74. const putRecords: (
  75. client: FirehoseClient,
  76. streamName: string,
  77. records: FirehoseRecord[],
  78. attempt?: number,
  79. ) => Effect.Effect<number, IngestError> = Effect.fn("Ingest.putRecords")(function* (
  80. client,
  81. streamName,
  82. records,
  83. attempt = 1,
  84. ) {
  85. const result = yield* Effect.tryPromise({
  86. try: () => client.send(new PutRecordBatchCommand({ DeliveryStreamName: streamName, Records: records })),
  87. catch: (cause) =>
  88. new IngestError({ message: "Failed to write lake records to Firehose", failed: records.length, cause }),
  89. }).pipe(
  90. Effect.tapError(() =>
  91. Effect.logWarning(`firehose batch write failed ${JSON.stringify({ records: records.length, attempt })}`),
  92. ),
  93. )
  94. const failed =
  95. result.RequestResponses?.flatMap((item, index) => {
  96. const record = records[index]
  97. if (!item.ErrorCode || !record) return []
  98. return [record]
  99. }) ?? []
  100. if (failed.length === 0) return 0
  101. if (attempt >= MAX_FIREHOSE_ATTEMPTS) {
  102. yield* Effect.logWarning(
  103. `firehose batch failed ${JSON.stringify({ records: failed.length, attempts: MAX_FIREHOSE_ATTEMPTS })}`,
  104. )
  105. return failed.length
  106. }
  107. yield* Effect.logWarning(
  108. `firehose batch retrying ${JSON.stringify({ records: failed.length, attempt: attempt + 1 })}`,
  109. )
  110. yield* Effect.sleep(`${250 * 2 ** (attempt - 1)} millis`)
  111. return yield* putRecords(client, streamName, failed, attempt + 1)
  112. })
  113. function countRoutedEvents(events: unknown[]) {
  114. let records = 0
  115. let unsupported = 0
  116. for (const event of events) {
  117. if (!isRecord(event)) continue
  118. if (routeEvent(event)) records++
  119. else unsupported++
  120. }
  121. return { records, unsupported }
  122. }
  123. function isRecord(item: unknown): item is IngestEvent {
  124. return Boolean(item) && typeof item === "object" && !Array.isArray(item)
  125. }
  126. function routeEvent(event: IngestEvent): LakeRoute | undefined {
  127. if (typeof event._datalake_key !== "string") return
  128. const match = event._datalake_key.match(LAKE_TYPE)
  129. if (!match?.[1] || !match[2]) return
  130. return {
  131. database: match[1],
  132. table: match[2],
  133. }
  134. }
  135. function toFirehoseRecord(event: IngestEvent, route: LakeRoute): FirehoseRecord {
  136. return {
  137. Data: Buffer.from(
  138. JSON.stringify({
  139. ...Object.fromEntries(Object.entries(event).filter(([key]) => key !== "_datalake_key")),
  140. _lake_database: route.database,
  141. _lake_table: route.table,
  142. _lake_operation: "insert" as const,
  143. }),
  144. ),
  145. }
  146. }