{"record":{"id":"34056671d7edd41b","repo":"hcengineering/platform","slug":"consumer-disconnected-from-queue","errorCode":null,"errorMessage":"consumer disconnected from queue","messagePattern":"consumer disconnected from queue","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"foundations/server/packages/kafka/src/index.ts","lineNumber":323,"sourceCode":"            await heartbeat()\n            await new Promise((resolve) => setTimeout(resolve, to * retryDelay))\n            if (to < maxRetryDelay) {\n              to++\n            }\n          }\n        }\n      }\n    })\n  }\n\n  async doConnect (): Promise<void> {\n    this.cc.on('consumer.connect', () => {\n      this.connected = true\n      this.ctx.info('consumer connected to queue')\n    })\n    this.cc.on('consumer.disconnect', () => {\n      this.connected = false\n      this.ctx.warn('consumer disconnected from queue')\n    })\n    await this.cc.connect()\n  }\n\n  async doSubscribe (): Promise<void> {\n    await this.cc.subscribe({\n      topic: getKafkaTopicId(this.topic, this.config),\n      fromBeginning: this.options?.fromBegining\n    })\n  }\n\n  isConnected (): boolean {\n    return this.connected\n  }\n\n  close (): Promise<void> {\n    return this.cc.disconnect()\n  }","sourceCodeStart":305,"sourceCodeEnd":341,"githubUrl":"https://github.com/hcengineering/platform/blob/63e28dc96483967b2fc21c881b3f1023c1de7718/foundations/server/packages/kafka/src/index.ts#L305-L341","documentation":"KafkaJS fires a consumer.disconnect event when the underlying consumer connection to the Kafka brokers drops. The kafka package's doConnect registers a handler that sets connected=false and logs this warning. Messages cannot be consumed until a reconnect succeeds; KafkaJS will normally auto-reconnect.","triggerScenarios":"During start() -> doConnect(), the consumer's broker connection terminates — broker restart, network partition, idle connection reaped, TLS/auth failure, or group rebalancing kicking the member out.","commonSituations":"Kafka broker restarts or rolling upgrades; network flakiness between server and Kafka cluster; connection_max_idle_ms expiring idle connections; misconfigured advertised listeners causing unreachable broker addresses.","solutions":["Check Kafka broker health and listener/advertised.listeners configuration","Verify network connectivity (host, port, TLS) between the server and brokers","Rely on KafkaJS auto-reconnect or add explicit reconnect/retry logic around consumer start","Review broker logs for group rebalance or auth failures that triggered the disconnect","Increase connectionTimeout/sessionTimeout values if timeouts cause the disconnect"],"exampleFix":"// before\nawait this.cc.connect()\n// after\nawait this.cc.connect()\nthis.cc.on('consumer.crash', async (e) => {\n  this.ctx.error('consumer crashed', { error: e })\n  await this.start(this.ctx) // re-establish connection/subscription\n})","handlingStrategy":"retry","validationCode":"const reachable = await net.connect({ host: kafkaHost, port: kafkaPort })\n  .then(s => { s.destroy(); return true })\n  .catch(() => false)\nif (!reachable) throw new Error('Kafka broker unreachable before consumer start')","typeGuard":"function isDisconnectEvent(e: unknown): e is { payload: { clientId: string } } {\n  return typeof e === 'object' && e !== null && 'payload' in e\n}","tryCatchPattern":"this.cc.on('consumer.disconnect', async () => {\n  this.connected = false\n  await backoffRetry(() => this.doConnect(), { retries: Infinity, baseMs: 1000 })\n})","preventionTips":["Use KafkaJS auto-restart/retry options and monitor consumer.disconnect events","Keep TCP keepalive on long-lived broker connections","Pin brokers behind stable DNS; verify advertised.listeners","Alert on prolonged connected=false periods"],"tags":["kafka","connection","reconnect"],"backgroundTag":"consumer-disconnected","analyzedSha":"63e28dc96483967b2fc21c881b3f1023c1de7718","analyzedAt":"2026-08-29T15:21:27.377Z","schemaVersion":2},"datasetVersion":"2026-08-29T17:17:51.833Z"}