Skip to content

Incorrect panic handling in session reduce #160

Description

@BulkBeing

I'm running a source->session reduce -> log sink pipeline using the Typescript SDK. When the exception thrown in Javascript is re-thrown in the session_reduce implementation, none of the containers becomes unhealthy or restarts.

ts-crash-session-reduce-0-v2ij0   3/3     Running     0          16m
Image

Numa logs

{"timestamp":"2025-12-17T11:46:31.490261Z","level":"ERROR","message":"Failed to send message reduce task, task aborted","e":"SendError { .. }","target":"numaflow_core::reduce::reducer::unaligned::reducer"}
{"timestamp":"2025-12-17T11:46:31.490265Z","level":"ERROR","message":"Failed to send message reduce task, task aborted","e":"SendError { .. }","target":"numaflow_core::reduce::reducer::unaligned::reducer"}
{"timestamp":"2025-12-17T11:46:32.189699Z","level":"INFO","message":"Processed messages per second","processed":"20","target":"numaflow_core::tracker"}
{"timestamp":"2025-12-17T11:46:32.492401Z","level":"ERROR","message":"Failed to send message reduce task, task aborted","e":"SendError { .. }","target":"numaflow_core::reduce::reducer::unaligned::reducer"}
{"timestamp":"2025-12-17T11:46:32.492544Z","level":"ERROR","message":"Failed to send message reduce task, task aborted","e":"SendError { .. }","target":"numaflow_core::reduce::reducer::unaligned::reducer"}

UDF logs:

{"level":30,"time":1765971869111,"pid":1,"hostname":"ts-crash-session-reduce-0","msg":"SessionReduceCounter: 180913 milliseconds since last crash, keys: 052b3034-38e2-462d-821f-0285fb152674"}
{"level":30,"time":1765971869114,"pid":1,"hostname":"ts-crash-session-reduce-0","msg":"3 Minutes since last crash. Simulating crash now"}
{"level":30,"time":1765971869121,"pid":1,"hostname":"ts-crash-session-reduce-0","msg":"SessionReduceCounter: 10 milliseconds since last crash, keys: 0e127135-f399-4a7e-bc14-2af063c31e8c"}
{"level":30,"time":1765971869121,"pid":1,"hostname":"ts-crash-session-reduce-0","msg":"0.01 seconds since last crash"}
{"level":30,"time":1765971869122,"pid":1,"hostname":"ts-crash-session-reduce-0","msg":"Counter: 7420"}
[ERROR] Error executing iterator returned by user-defined session reduce function: Error { status: "GenericFailure", reason: "Error: Simulated crash" }
{"level":30,"time":1765971869122,"pid":1,"hostname":"ts-crash-session-reduce-0","msg":"Counter: 7420"}
{"level":30,"time":1765971869123,"pid":1,"hostname":"ts-crash-session-reduce-0","msg":"Counter: 7420"}
{"level":30,"time":1765971869124,"pid":1,"hostname":"ts-crash-session-reduce-0","msg":"Counter: 7420"}
{"level":30,"time":1765971869124,"pid":1,"hostname":"ts-crash-session-reduce-0","msg":"Counter: 7420"}
{"level":30,"time":1765971869125,"pid":1,"hostname":"ts-crash-session-reduce-0","msg":"Counter: 7420"}

UDF implementation:

class SessionReduceCounter implements sessionReduce.SessionReducer {
    counter: number = 0
    async *sessionReduceFn(
        keys: string[],
        datums: AsyncIterableIterator<sessionReduce.Datum>
    ): AsyncIterableIterator<sessionReduce.Message> {
        const now = Date.now()
        logger.info(
            `SessionReduceCounter: ${
                (now - lastCrashTime) / 1000
            } seconds since last crash, keys: ${keys.join(',')}`
        )
        if (lastCrashTime === 0) {
            // Initialize the timer on the first request but do not crash yet.
            lastCrashTime = now
        } else if (now - lastCrashTime >= CRASH_INTERVAL) {
            logger.info('3 Minutes since last crash. Simulating crash now')
            // Update the last crash time so that, even with many concurrent requests,
            // only a single request in each interval will see the simulated crash.
            lastCrashTime = now
            throw new Error('Simulated crash')
        } else {
            logger.info(
                `${(now - lastCrashTime) / 1000} seconds since last crash`
            )
        }

        for await (const _ of datums) {
            this.counter += 1
        }
        logger.info(`Counter: ${this.counter}`)

        yield {
            keys,
            value: Buffer.from(this.counter.toString()),
        }
    }

    async accumulatorFn(): Promise<Buffer> {
        return Buffer.from(this.counter.toString())
    }

    async mergeAccumulatorFn(accumulator: Buffer): Promise<void> {
        this.counter += accumulator.readUInt8(0)
    }
}

Reduce vertex spec:

    - name: session-reduce
      scale:
        min: 1
        max: 1
      udf:
        groupBy:
          window:
            session:
              timeout: 10s
          keyed: true
          storage:
            persistentVolumeClaim:
              volumeSize: 1Gi
        container:
          image: js-endurance-session-reduce-crash:stable
          imagePullPolicy: Never

Source vertex sends message that looks like:

{
    payload: Buffer.from(JSON.stringify(payload), 'utf-8'),
    offset: offset,
    eventTime: new Date(),
    keys: [crypto.randomUUID()],
    headers: { source: 'ad-events' },
    userMetadata: userMetadata,
}

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions