@@ -3,9 +3,10 @@ import {
33 type EventsQueuePayloadIncomingEvent ,
44 KAFKA_EVENTS_TOPIC ,
55 KAFKA_PARTITIONS_CONCURRENT ,
6+ type KafkaEventsMessage ,
67 kafkaLogger ,
7- type KafkaMessage ,
88} from '@openpanel/queue' ;
9+ import { forEach } from 'hwp' ;
910import { logger } from '../utils/logger' ;
1011import { markEventsActivity } from '../utils/worker-heartbeat' ;
1112import { incomingEvent } from './events.incoming-event' ;
@@ -14,103 +15,40 @@ export interface KafkaConsumerHandle {
1415 stop : ( ) => Promise < void > ;
1516}
1617
17- // Heartbeat every N messages within a per-key group. The default kafkajs
18- // sessionTimeout is 30s — calling heartbeat every 16 messages keeps us
19- // comfortably under that even for slow handlers.
20- const HEARTBEAT_EVERY = 16 ;
21-
2218export async function startKafkaEventsConsumer ( ) : Promise < KafkaConsumerHandle > {
2319 const consumer = createKafkaEventsConsumer ( ) ;
24- await consumer . connect ( ) ;
25- await consumer . subscribe ( {
26- topic : KAFKA_EVENTS_TOPIC ,
27- fromBeginning : false ,
28- } ) ;
2920
30- consumer . on ( consumer . events . HEARTBEAT , markEventsActivity ) ;
21+ consumer . on ( ' consumer:heartbeat:end' , markEventsActivity ) ;
3122
32- await consumer . run ( {
33- partitionsConsumedConcurrently : KAFKA_PARTITIONS_CONCURRENT ,
34- eachBatchAutoResolve : false ,
35- eachBatch : async ( {
36- batch,
37- resolveOffset,
38- heartbeat,
39- isRunning,
40- isStale,
41- } ) => {
42- if ( batch . messages . length === 0 ) {
43- return ;
44- }
23+ const stream = await consumer . consume ( { topics : [ KAFKA_EVENTS_TOPIC ] } ) ;
4524
46- // Group by partition key (= deviceId or `${projectId}:${profileId}`).
47- // Same-key messages stay serial so sessionBuffer/session-end-job state
48- // can't race; different keys run in parallel via Promise.all.
49- // Keyless messages get their own singleton group.
50- const groups = new Map < string , KafkaMessage [ ] > ( ) ;
51- for ( const m of batch . messages ) {
52- const key = m . key ? m . key . toString ( ) : `__no_key__:${ m . offset } ` ;
53- const arr = groups . get ( key ) ;
54- if ( arr ) {
55- arr . push ( m ) ;
56- } else {
57- groups . set ( key , [ m ] ) ;
25+ // Same-key messages stay serial; different keys parallelize up to
26+ // KAFKA_PARTITIONS_CONCURRENT. The key is deviceId or projectId:profileId
27+ // (set in track.controller.ts), so this preserves per-session ordering
28+ // for the sessionBuffer / session-end pipeline.
29+ const tails = new Map < string , Promise < void > > ( ) ;
30+
31+ const consumeLoop = forEach (
32+ stream ,
33+ async ( message : KafkaEventsMessage ) => {
34+ const key =
35+ message . key && message . key . length > 0
36+ ? message . key . toString ( )
37+ : `__no_key__:${ message . offset . toString ( ) } ` ;
38+ const prev = tails . get ( key ) ?? Promise . resolve ( ) ;
39+ const next = prev . then ( ( ) => processMessage ( message ) ) ;
40+ tails . set ( key , next ) ;
41+ try {
42+ await next ;
43+ } finally {
44+ if ( tails . get ( key ) === next ) {
45+ tails . delete ( key ) ;
5846 }
5947 }
60-
61- await Promise . all (
62- [ ...groups . values ( ) ] . map ( async ( msgs ) => {
63- let processed = 0 ;
64- for ( const m of msgs ) {
65- if ( ! isRunning ( ) || isStale ( ) ) {
66- return ;
67- }
68-
69- if ( m . value ) {
70- let payload : EventsQueuePayloadIncomingEvent [ 'payload' ] | null =
71- null ;
72- try {
73- payload = JSON . parse (
74- m . value . toString ( )
75- ) as EventsQueuePayloadIncomingEvent [ 'payload' ] ;
76- } catch ( err ) {
77- logger . error (
78- { err, partition : batch . partition , offset : m . offset } ,
79- 'kafka message parse failed'
80- ) ;
81- }
82- if ( payload ) {
83- try {
84- await incomingEvent ( payload ) ;
85- } catch ( err ) {
86- // Match the previous eachMessage behaviour: log and ack.
87- // At-most-once on handler exceptions; failures here would
88- // otherwise block the partition.
89- logger . error (
90- {
91- err,
92- partition : batch . partition ,
93- offset : m . offset ,
94- projectId : payload . projectId ,
95- } ,
96- 'kafka incomingEvent handler failed'
97- ) ;
98- }
99- }
100- }
101-
102- resolveOffset ( m . offset ) ;
103- processed += 1 ;
104- if ( processed % HEARTBEAT_EVERY === 0 ) {
105- await heartbeat ( ) ;
106- }
107- }
108- } )
109- ) ;
110-
111- await heartbeat ( ) ;
112- markEventsActivity ( ) ;
11348 } ,
49+ KAFKA_PARTITIONS_CONCURRENT
50+ ) . catch ( ( err : unknown ) => {
51+ logger . error ( { err } , 'kafka consumer stream errored' ) ;
11452 } ) ;
11553
11654 kafkaLogger . info (
@@ -120,10 +58,50 @@ export async function startKafkaEventsConsumer(): Promise<KafkaConsumerHandle> {
12058 } ,
12159 'kafka events consumer running'
12260 ) ;
61+ markEventsActivity ( ) ;
12362
12463 return {
12564 stop : async ( ) => {
126- await consumer . disconnect ( ) ;
65+ await stream . close ( ) ;
66+ await consumeLoop ;
67+ await consumer . close ( ) ;
12768 } ,
12869 } ;
12970}
71+
72+ async function processMessage ( message : KafkaEventsMessage ) : Promise < void > {
73+ if ( ! message . value || message . value . length === 0 ) {
74+ return ;
75+ }
76+ let payload : EventsQueuePayloadIncomingEvent [ 'payload' ] | null = null ;
77+ try {
78+ payload = JSON . parse (
79+ message . value . toString ( )
80+ ) as EventsQueuePayloadIncomingEvent [ 'payload' ] ;
81+ } catch ( err ) {
82+ logger . error (
83+ {
84+ err,
85+ partition : message . partition ,
86+ offset : message . offset . toString ( ) ,
87+ } ,
88+ 'kafka message parse failed'
89+ ) ;
90+ return ;
91+ }
92+ try {
93+ await incomingEvent ( payload ) ;
94+ } catch ( err ) {
95+ // At-most-once on handler exceptions: log and let autocommit advance.
96+ // Throwing here would block the partition until human intervention.
97+ logger . error (
98+ {
99+ err,
100+ partition : message . partition ,
101+ offset : message . offset . toString ( ) ,
102+ projectId : payload . projectId ,
103+ } ,
104+ 'kafka incomingEvent handler failed'
105+ ) ;
106+ }
107+ }
0 commit comments