| 230 | }; |
| 231 | |
| 232 | export const disconnectKafka = async (): Promise<void> => { |
| 233 | const tasks: Promise<unknown>[] = []; |
| 234 | for (const c of consumers) { |
| 235 | tasks.push( |
| 236 | c.disconnect().catch((err) => { |
| 237 | kafkaLogger.error({ err }, 'kafka consumer disconnect error'); |
| 238 | }) |
| 239 | ); |
| 240 | } |
| 241 | consumers.clear(); |
| 242 | if (producer) { |
| 243 | const p = producer; |
| 244 | producer = null; |
| 245 | producerConnectPromise = null; |
| 246 | tasks.push( |
| 247 | p.disconnect().catch((err) => { |
| 248 | kafkaLogger.error({ err }, 'kafka producer disconnect error'); |
| 249 | }) |
| 250 | ); |
| 251 | } |
| 252 | await Promise.all(tasks); |
| 253 | }; |