{"_id":"@abinashpatri/kafka","_rev":"2-8d5d05a9723ab6f052f290704efe76fc","name":"@abinashpatri/kafka","dist-tags":{"latest":"1.0.1"},"versions":{"1.0.0":{"name":"@abinashpatri/kafka","version":"1.0.0","keywords":["events","kafka","messaging","typescript"],"author":{"name":"Abinash Patri"},"license":"MIT","_id":"@abinashpatri/kafka@1.0.0","maintainers":[{"name":"abinashpatri","email":"abinashpatri33@gmail.com"}],"dist":{"shasum":"d0d0553b861e8cfe64f101db4fcafb4a049309ee","tarball":"https://registry.npmjs.org/@abinashpatri/kafka/-/kafka-1.0.0.tgz","fileCount":17,"integrity":"sha512-EkitgIyM5Ul0BSk8G6Wx0gPbCRNuQz148rP6sK6BZRCZuGwKTvNulhM6NISksVD8Cj1vILqxFvyX3j+aI14IHQ==","signatures":[{"sig":"MEUCIE90JPsQ2+HGxozYvofST14QORke+O69ZhvxbVAQU0HxAiEA11kb+yTfyOdexzRbUOdDZ0gk1D5hI609yEs6B/ZKXNk=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":78039},"main":"./dist/index.js","type":"commonjs","types":"./dist/index.d.ts","module":"./dist/index.mjs","engines":{"node":">=18"},"exports":{".":{"types":"./dist/index.d.ts","import":"./dist/index.mjs","require":"./dist/index.js"},"./kafka":{"types":"./dist/kafka/index.d.ts","import":"./dist/kafka/index.mjs","require":"./dist/kafka/index.js"}},"scripts":{"dev":"tsup --watch","build":"tsup","typecheck":"tsc --noEmit","prepublishOnly":"npm run typecheck && npm run build"},"_npmUser":{"name":"abinashpatri","email":"abinashpatri33@gmail.com"},"_npmVersion":"11.6.2","description":"Production-grade Kafka event utility library","directories":{},"sideEffects":false,"_nodeVersion":"24.11.1","dependencies":{"kafkajs":"^2.2.4"},"_hasShrinkwrap":false,"devDependencies":{"tsup":"^8.5.1","eslint":"^10.1.0","prettier":"^3.8.1","typescript":"^5.9.3","@types/node":"^25.5.0"},"_npmOperationalInternal":{"tmp":"tmp/kafka_1.0.0_1774250484264_0.9999674431911085","host":"s3://npm-registry-packages-npm-production"}},"1.0.1":{"name":"@abinashpatri/kafka","version":"1.0.1","description":"Production-grade Kafka event utility library","main":"./dist/index.js","module":"./dist/index.mjs","types":"./dist/index.d.ts","exports":{".":{"types":"./dist/index.d.ts","require":"./dist/index.js","import":"./dist/index.mjs"},"./kafka":{"types":"./dist/kafka/index.d.ts","require":"./dist/kafka/index.js","import":"./dist/kafka/index.mjs"}},"sideEffects":false,"scripts":{"build":"tsup","dev":"tsup --watch","typecheck":"tsc --noEmit","prepublishOnly":"npm run typecheck && npm run build"},"keywords":["events","kafka","messaging","typescript"],"author":{"name":"Abinash Patri"},"license":"MIT","type":"commonjs","engines":{"node":">=18"},"devDependencies":{"@types/node":"^25.5.0","eslint":"^10.1.0","prettier":"^3.8.1","tsup":"^8.5.1","typescript":"^5.9.3"},"dependencies":{"kafkajs":"^2.2.4"},"_id":"@abinashpatri/kafka@1.0.1","_nodeVersion":"24.11.1","_npmVersion":"11.6.2","dist":{"integrity":"sha512-WBCaSi/ZiA4/09Qs3TPcMbPxm85oOgKgsTLk7uRFInGO1iwTyAhmIW0+JVzj1AkaoXlwTfwgyR5CkWyLEw1erA==","shasum":"84d3559ceb83c40b5e8f7cdba4e8654934229eca","tarball":"https://registry.npmjs.org/@abinashpatri/kafka/-/kafka-1.0.1.tgz","fileCount":17,"unpackedSize":77962,"signatures":[{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEUCIBBD2RmuLLEQTL7nfK/8sb+rOjIan54+2aaCRSRXEgd+AiEA+S8OvtTrYSpOH4vvn9rPrym7hk1Ap2rRnMUdsBFYAuA="}]},"_npmUser":{"name":"abinashpatri","email":"abinashpatri33@gmail.com"},"directories":{},"maintainers":[{"name":"abinashpatri","email":"abinashpatri33@gmail.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/kafka_1.0.1_1774845505766_0.775646046333913"},"_hasShrinkwrap":false}},"time":{"created":"2026-03-23T07:21:24.111Z","modified":"2026-03-30T04:38:26.004Z","1.0.0":"2026-03-23T07:21:24.409Z","1.0.1":"2026-03-30T04:38:25.904Z"},"author":{"name":"Abinash Patri"},"license":"MIT","keywords":["events","kafka","messaging","typescript"],"description":"Production-grade Kafka event utility library","maintainers":[{"name":"abinashpatri","email":"abinashpatri33@gmail.com"}],"readme":"# @abinashpatri/kafka\n\nProduction-grade Kafka eventing helpers with typed APIs, retry/DLQ flows, and scoped clients for multi-service runtimes.\n\n## Install\n\n```bash\nnpm install @abinashpatri/kafka\n```\n\n## Compatibility\n\n- Node.js: `>=18`\n- Module formats: CommonJS and ESM\n- Type support: bundled `.d.ts`\n\n## Import Patterns\n\nUse either root namespace imports or subpath imports.\n\n```ts\n// Root namespace import\nimport { kafka } from \"@abinashpatri/kafka\";\n\n// Subpath import\nimport * as kafkaApi from \"@abinashpatri/kafka/kafka\";\n```\n\n## Feature Overview\n\n- Connect/disconnect producer client.\n- Typed publish and consume helpers.\n- Retry topic helpers (`<topic>.retry.<n>`).\n- DLQ helper (`<topic>.dlq`).\n- Consumer lifecycle control (`stop`, `disconnect`).\n- Retry backoff + jitter controls.\n- Scoped client factory for non-shared state.\n\n## Quick Start (TypeScript)\n\n```ts\nimport { kafka } from \"@abinashpatri/kafka\";\n\ntype UserCreatedEvent = {\n  eventId: string;\n  userId: string;\n  email: string;\n  createdAt: string;\n};\n\nawait kafka.connect({\n  clientId: \"user-service\",\n  brokers: [\"localhost:9092\"],\n});\n\nawait kafka.publish<UserCreatedEvent>({\n  topic: \"user.created\",\n  key: \"user_123\",\n  message: {\n    eventId: \"evt_1\",\n    userId: \"user_123\",\n    email: \"user@example.com\",\n    createdAt: new Date().toISOString(),\n  },\n});\n\nconst kafkaConsumer = await kafka.consume<UserCreatedEvent>({\n  topic: \"user.created\",\n  groupId: \"analytics-group\",\n  retryLimit: 5,\n  retryBaseDelayMs: 500,\n  retryBackoffMultiplier: 2,\n  retryMaxDelayMs: 30000,\n  retryJitterMs: 250,\n  handler: async (event) => {\n    console.log(\"analytics processing\", event.userId);\n  },\n});\n```\n\n## Quick Start (JavaScript)\n\n```js\nconst { kafka } = require(\"@abinashpatri/kafka\");\n\nasync function runKafka() {\n  await kafka.connect({\n    clientId: \"order-service\",\n    brokers: [\"localhost:9092\"],\n  });\n\n  await kafka.publish({\n    topic: \"order.created\",\n    key: \"order_1\",\n    message: {\n      eventId: \"evt_100\",\n      orderId: \"order_1\",\n      total: 129.99,\n      createdAt: new Date().toISOString(),\n    },\n  });\n\n  return kafka.consume({\n    topic: \"order.created\",\n    groupId: \"fulfillment-group\",\n    retryLimit: 5,\n    retryBaseDelayMs: 500,\n    retryBackoffMultiplier: 2,\n    retryMaxDelayMs: 30000,\n    retryJitterMs: 250,\n    handler: async (event) => {\n      console.log(\"fulfillment:\", event.orderId);\n    },\n  });\n}\n\nrunKafka().catch(console.error);\n```\n\n## API Reference\n\n### Root export\n\n- `kafka` namespace\n\n### Kafka API\n\n- `kafka.connect({ clientId, brokers })`\n- `kafka.disconnect()`\n- `kafka.publish({ topic, message, key?, headers?, client? })`\n- `kafka.consume({ topic, groupId, handler, retryLimit?, retryBaseDelayMs?, retryBackoffMultiplier?, retryMaxDelayMs?, retryJitterMs?, client? })`\n- `kafka.createKafkaClient()`\n- `kafka.createScopedKafkaClient()`\n- `kafka.sendToDLQ(baseTopic, message)`\n- `kafka.getRetryTopic(baseTopic, retryCount)`\n- `kafka.getDLQTopic(baseTopic)`\n- `kafka.getRetryCount(headers)`\n- `kafka.buildRetryHeaders(retryCount)`\n- `kafka.isProcessed(id)` / `kafka.markProcessed(id)` / `kafka.resetProcessedState()`\n\nKafka consumer control returned by `consume()`:\n\n- `stop()`\n- `disconnect()`\n\n## Reliability Semantics\n\n- Delivery model is generally **at-least-once**.\n- Handlers should be idempotent.\n- Retries are redirected to retry topics, then DLQ after retry limit.\n\n## Graceful Shutdown\n\nUse returned consumer controls and connection disconnect on process shutdown.\n\n```ts\nimport { kafka } from \"@abinashpatri/kafka\";\n\nconst kafkaControl = await kafka.consume({\n  topic: \"user.created\",\n  groupId: \"worker\",\n  handler: async () => {},\n});\n\nprocess.on(\"SIGTERM\", async () => {\n  await Promise.all([kafkaControl.disconnect(), kafka.disconnect()]);\n  process.exit(0);\n});\n```\n\n## License\n\nMIT License.\n","readmeFilename":"README.md"}