{"_id":"@anyq/kafka","_rev":"7-50f333807f0642c1151607f0c6aabb9a","name":"@anyq/kafka","dist-tags":{"latest":"0.5.0"},"versions":{"0.1.0":{"name":"@anyq/kafka","version":"0.1.0","keywords":["queue","kafka","message-queue","streaming","anyq"],"author":{"name":"Shantanu Sharma"},"license":"MIT","_id":"@anyq/kafka@0.1.0","maintainers":[{"name":"sns45","email":"notifyshantanu@gmail.com"}],"homepage":"https://github.com/sns45/anyq#readme","bugs":{"url":"https://github.com/sns45/anyq/issues"},"dist":{"shasum":"48c1ba01f70f2e94ada1c17b11f9d3fe79f25674","tarball":"https://registry.npmjs.org/@anyq/kafka/-/kafka-0.1.0.tgz","fileCount":10,"integrity":"sha512-RKTdrqdEhCqCdVL2SmUhHLd/d5Ggc8fFecRfSk46e3KdhuoBW4s7RENe9rOEF8htRL3bg850JL6IqI1qBnCdmg==","signatures":[{"sig":"MEUCIQD3HPS4ZopL1DUb+heZqlph+nyA1UT/A2KyYKqR/P8nFgIgIn1BpEDxDpSrxONUAAum+lh6JrAI1er2jsrlcVRxpcs=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":587239},"main":"./dist/index.js","type":"module","types":"./dist/index.d.ts","exports":{".":{"types":"./dist/index.d.ts","import":"./dist/index.js","require":"./dist/index.js"}},"gitHead":"1210b51fe4a55144ebf5bc798db0732b4c89ddba","scripts":{"test":"bun test","build":"bun build ./src/index.ts --outdir ./dist --target bun && bun run build:types","clean":"rm -rf dist","typecheck":"tsc --noEmit","build:types":"tsc --emitDeclarationOnly --outDir dist"},"_npmUser":{"name":"sns45","email":"notifyshantanu@gmail.com"},"repository":{"url":"git+https://github.com/sns45/anyq.git","type":"git","directory":"packages/kafka"},"_npmVersion":"10.9.2","description":"Apache Kafka adapter for anyq","directories":{},"_nodeVersion":"22.17.0","dependencies":{"kafkajs":"^2.2.4","@anyq/core":"workspace:*"},"publishConfig":{"access":"public"},"_hasShrinkwrap":false,"devDependencies":{"@types/bun":"^1.1.0","typescript":"^5.0.0"},"peerDependencies":{"@anyq/core":"workspace:*"},"_npmOperationalInternal":{"tmp":"tmp/kafka_0.1.0_1765990968961_0.9009351789441264","host":"s3://npm-registry-packages-npm-production"}},"0.1.1":{"name":"@anyq/kafka","version":"0.1.1","keywords":["queue","kafka","message-queue","streaming","anyq"],"author":{"name":"Shantanu Sharma"},"license":"MIT","_id":"@anyq/kafka@0.1.1","maintainers":[{"name":"sns45","email":"notifyshantanu@gmail.com"}],"homepage":"https://github.com/sns45/anyq#readme","bugs":{"url":"https://github.com/sns45/anyq/issues"},"dist":{"shasum":"092bcda5f5f869d31acc86a58d2ccf5996a70914","tarball":"https://registry.npmjs.org/@anyq/kafka/-/kafka-0.1.1.tgz","fileCount":11,"integrity":"sha512-80eBNawnHXxflXivBpYEnKaij9vImPapbUBflzp61DhIj0AWRVo5+5RznJF5N00Ap0NR9/dr46JCNnx78okzlQ==","signatures":[{"sig":"MEYCIQDl5X1/SORc9mRtw2iVicfiqV8jWkB1/yjG72eCcEyS/gIhAPIkcRYnktAQ9GNaOtdxkNBe1/iFsKWLZxDORwem4HUD","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":589285},"main":"./dist/index.js","type":"module","types":"./dist/index.d.ts","exports":{".":{"types":"./dist/index.d.ts","import":"./dist/index.js","require":"./dist/index.js"}},"gitHead":"30128d564178a5a1d25b736b605bf69644178a0f","scripts":{"test":"bun test","build":"bun build ./src/index.ts --outdir ./dist --target bun && bun run build:types","clean":"rm -rf dist","typecheck":"tsc --noEmit","build:types":"tsc --emitDeclarationOnly --outDir dist"},"_npmUser":{"name":"sns45","email":"notifyshantanu@gmail.com"},"repository":{"url":"git+https://github.com/sns45/anyq.git","type":"git","directory":"packages/kafka"},"_npmVersion":"10.9.2","description":"Apache Kafka adapter for anyq","directories":{},"_nodeVersion":"22.17.0","dependencies":{"kafkajs":"^2.2.4","@anyq/core":"workspace:*"},"publishConfig":{"access":"public"},"_hasShrinkwrap":false,"devDependencies":{"@types/bun":"^1.1.0","typescript":"^5.0.0"},"peerDependencies":{"@anyq/core":"workspace:*"},"_npmOperationalInternal":{"tmp":"tmp/kafka_0.1.1_1765991382912_0.6518730917049407","host":"s3://npm-registry-packages-npm-production"}},"0.3.0":{"name":"@anyq/kafka","version":"0.3.0","keywords":["queue","kafka","message-queue","streaming","anyq"],"author":"Shantanu Sharma","license":"MIT","_id":"@anyq/kafka@0.3.0","maintainers":[{"name":"sns45","email":"notifyshantanu@gmail.com"}],"homepage":"https://github.com/sns45/anyq#readme","bugs":{"url":"https://github.com/sns45/anyq/issues"},"dist":{"shasum":"8a5d8f69e0db58c77989d4351edfe094c4278ea3","tarball":"https://registry.npmjs.org/@anyq/kafka/-/kafka-0.3.0.tgz","fileCount":11,"integrity":"sha512-74BTF6CuYgvx+NPBUlfrs1TTXtl6rHYi5dyefM/BEyyg1ACwznXtD8nO7PB22SRyo2/KFkSBCr/JbIq2796sjw==","signatures":[{"sig":"MEQCIG0VHASbFs8ccutNBB8jH699spIj+9rDWT7JayXZkPYVAiAWfxN5kLbTxKf2WD+wTwg3mUF1nDEX9xkoa+oLuPaBvg==","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":596776},"main":"./dist/index.js","type":"module","types":"./dist/index.d.ts","shasum":"8a5d8f69e0db58c77989d4351edfe094c4278ea3","exports":{".":{"types":"./dist/index.d.ts","import":"./dist/index.js","require":"./dist/index.js"}},"scripts":{"test":"bun test","build":"bun build ./src/index.ts --outdir ./dist --target bun && bun run build:types","clean":"rm -rf dist","typecheck":"tsc --noEmit","build:types":"tsc --emitDeclarationOnly --outDir dist"},"_npmUser":{"name":"sns45","email":"notifyshantanu@gmail.com"},"_integrity":"sha512-74BTF6CuYgvx+NPBUlfrs1TTXtl6rHYi5dyefM/BEyyg1ACwznXtD8nO7PB22SRyo2/KFkSBCr/JbIq2796sjw==","repository":{"url":"https://github.com/sns45/anyq","type":"git","directory":"packages/kafka"},"_npmVersion":"10.8.3","description":"Apache Kafka adapter for anyq","directories":{},"_nodeVersion":"24.3.0","dependencies":{"kafkajs":"^2.2.4","@anyq/core":"0.3.0"},"publishConfig":{"access":"public"},"_hasShrinkwrap":false,"devDependencies":{"@types/bun":"^1.1.0","typescript":"^5.0.0"},"peerDependencies":{"@anyq/core":"0.3.0"},"_npmOperationalInternal":{"tmp":"tmp/kafka_0.3.0_1781961979316_0.31443947589844856","host":"s3://npm-registry-packages-npm-production"}},"0.3.1":{"name":"@anyq/kafka","version":"0.3.1","keywords":["queue","kafka","message-queue","streaming","anyq"],"author":"Shantanu Sharma","license":"MIT","_id":"@anyq/kafka@0.3.1","maintainers":[{"name":"sns45","email":"notifyshantanu@gmail.com"}],"homepage":"https://github.com/sns45/anyq#readme","bugs":{"url":"https://github.com/sns45/anyq/issues"},"dist":{"shasum":"4ee66cdef692c73833be147b2749c47f0dfa83b8","tarball":"https://registry.npmjs.org/@anyq/kafka/-/kafka-0.3.1.tgz","fileCount":11,"integrity":"sha512-4O6OaisVfVqdRZ7vJxuDPRP+jivQT9lKCi9HPNDDdEzBH5cQ0Rm/639y9bOmVSinOzzfbc5zhbtlqwiGeFJf+A==","signatures":[{"sig":"MEUCIE6PKnpwlvelzD07/iVqPLab69QWW84rYb0mGbE77X0ZAiEAh19P/BuJRB008cjIddf71KaBJgNix2TKf4F0MQNwQ8M=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":599514},"main":"./dist/index.js","type":"module","types":"./dist/index.d.ts","shasum":"4ee66cdef692c73833be147b2749c47f0dfa83b8","exports":{".":{"types":"./dist/index.d.ts","import":"./dist/index.js","require":"./dist/index.js"}},"scripts":{"test":"bun test","build":"bun build ./src/index.ts --outdir ./dist --target bun && bun run build:types","clean":"rm -rf dist","typecheck":"tsc --noEmit","build:types":"tsc --emitDeclarationOnly --outDir dist"},"_npmUser":{"name":"sns45","email":"notifyshantanu@gmail.com"},"_integrity":"sha512-4O6OaisVfVqdRZ7vJxuDPRP+jivQT9lKCi9HPNDDdEzBH5cQ0Rm/639y9bOmVSinOzzfbc5zhbtlqwiGeFJf+A==","repository":{"url":"https://github.com/sns45/anyq","type":"git","directory":"packages/kafka"},"_npmVersion":"10.8.3","description":"Apache Kafka adapter for anyq","directories":{},"_nodeVersion":"24.3.0","dependencies":{"kafkajs":"^2.2.4","@anyq/core":"0.3.0"},"publishConfig":{"access":"public"},"_hasShrinkwrap":false,"devDependencies":{"@types/bun":"^1.1.0","typescript":"^5.0.0"},"peerDependencies":{"@anyq/core":"0.3.0"},"_npmOperationalInternal":{"tmp":"tmp/kafka_0.3.1_1782054603052_0.4114394105301131","host":"s3://npm-registry-packages-npm-production"}},"0.3.2":{"name":"@anyq/kafka","version":"0.3.2","keywords":["queue","kafka","message-queue","streaming","anyq"],"author":"Shantanu Sharma","license":"MIT","_id":"@anyq/kafka@0.3.2","maintainers":[{"name":"sns45","email":"notifyshantanu@gmail.com"}],"homepage":"https://github.com/sns45/anyq#readme","bugs":{"url":"https://github.com/sns45/anyq/issues"},"dist":{"shasum":"a008ec2504ef2b47631b563924282175b3644094","tarball":"https://registry.npmjs.org/@anyq/kafka/-/kafka-0.3.2.tgz","fileCount":11,"integrity":"sha512-X93yHMSAIvuPVfLDcryeAs+M7eJI0BsbQ9Cvw6oJtf0we06eZH0IMUjGxTMCsXY3FUWyQ7Cu/6M7fI6Tmcf9zw==","signatures":[{"sig":"MEYCIQDj6ACTY8AkX4I6zPx0RjtWgjnvwz94RHst0w3TR04CRwIhALarAYTlWs+VhNY54hC401bsUs6HZH8PyidBck6IGCbL","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":599514},"main":"./dist/index.js","type":"module","types":"./dist/index.d.ts","shasum":"a008ec2504ef2b47631b563924282175b3644094","exports":{".":{"types":"./dist/index.d.ts","import":"./dist/index.js","require":"./dist/index.js"}},"scripts":{"test":"bun test","build":"bun build ./src/index.ts --outdir ./dist --target bun && bun run build:types","clean":"rm -rf dist","typecheck":"tsc --noEmit","build:types":"tsc --emitDeclarationOnly --outDir dist"},"_npmUser":{"name":"sns45","email":"notifyshantanu@gmail.com"},"_integrity":"sha512-X93yHMSAIvuPVfLDcryeAs+M7eJI0BsbQ9Cvw6oJtf0we06eZH0IMUjGxTMCsXY3FUWyQ7Cu/6M7fI6Tmcf9zw==","repository":{"url":"https://github.com/sns45/anyq","type":"git","directory":"packages/kafka"},"_npmVersion":"10.8.3","description":"Apache Kafka adapter for anyq","directories":{},"_nodeVersion":"24.3.0","dependencies":{"kafkajs":"^2.2.4","@anyq/core":"0.3.2"},"publishConfig":{"access":"public"},"_hasShrinkwrap":false,"devDependencies":{"@types/bun":"^1.1.0","typescript":"^5.0.0"},"peerDependencies":{"@anyq/core":"0.3.2"},"_npmOperationalInternal":{"tmp":"tmp/kafka_0.3.2_1784910270455_0.20367112810303256","host":"s3://npm-registry-packages-npm-production"}},"0.4.0":{"name":"@anyq/kafka","version":"0.4.0","keywords":["queue","kafka","message-queue","streaming","anyq"],"author":"Shantanu Sharma","license":"Apache-2.0","_id":"@anyq/kafka@0.4.0","maintainers":[{"name":"sns45","email":"notifyshantanu@gmail.com"}],"homepage":"https://github.com/sns45/anyq#readme","bugs":{"url":"https://github.com/sns45/anyq/issues"},"dist":{"shasum":"06afefb5b78d1a7263ac376f79a9b2c6b01bb4aa","tarball":"https://registry.npmjs.org/@anyq/kafka/-/kafka-0.4.0.tgz","fileCount":13,"integrity":"sha512-WH+XwUhYTb+TADaQtQ2uzWOMwZ41P7dqlFnqz8NxBIGw15lpE3PXPVtXy4G/33l7osjljumc6OWqoXnvrWR0QQ==","signatures":[{"sig":"MEYCIQD/MB7F5zjHz5wsVX0WNkvwcFUJim+95e56cFzTCGsgrAIhANfaaaTAB7uMTZ7CR1Mawi/A/+AFyRiR4BNSEWIHJqN2","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":613532},"main":"./dist/index.js","type":"module","types":"./dist/index.d.ts","shasum":"06afefb5b78d1a7263ac376f79a9b2c6b01bb4aa","exports":{".":{"types":"./dist/index.d.ts","import":"./dist/index.js","require":"./dist/index.js"}},"scripts":{"test":"bun test","build":"bun build ./src/index.ts --outdir ./dist --target bun && bun run build:types","clean":"rm -rf dist","typecheck":"tsc --noEmit","build:types":"tsc --emitDeclarationOnly --outDir dist"},"_npmUser":{"name":"sns45","email":"notifyshantanu@gmail.com"},"_integrity":"sha512-WH+XwUhYTb+TADaQtQ2uzWOMwZ41P7dqlFnqz8NxBIGw15lpE3PXPVtXy4G/33l7osjljumc6OWqoXnvrWR0QQ==","repository":{"url":"https://github.com/sns45/anyq","type":"git","directory":"packages/kafka"},"_npmVersion":"10.8.3","description":"Apache Kafka adapter for anyq","directories":{},"_nodeVersion":"26.3.0","dependencies":{"kafkajs":"^2.2.4","@anyq/core":"0.4.0"},"publishConfig":{"access":"public"},"_hasShrinkwrap":false,"devDependencies":{"@types/bun":"^1.1.0","typescript":"^5.0.0"},"peerDependencies":{"@anyq/core":"0.4.0"},"_npmOperationalInternal":{"tmp":"tmp/kafka_0.4.0_1787860574490_0.2264546251852111","host":"s3://npm-registry-packages-npm-production"}},"0.5.0":{"_id":"@anyq/kafka@0.5.0","bugs":{"url":"https://github.com/sns45/anyq/issues"},"dist":{"shasum":"186e7d558e69504fa97ec3cc7ed02c109f55723c","tarball":"https://registry.npmjs.org/@anyq/kafka/-/kafka-0.5.0.tgz","fileCount":13,"integrity":"sha512-6RYjuIkhlHMk8qxSWKrsinGRgEvFZPVrKrWXBze8QcCzev3Ot+8qxT1Z4oVcgYzPoxQY+1mF20OdRT90TZ4Xkw==","signatures":[{"sig":"MEQCICckgEoerjzp5d1BYqrCsP3vBFm2A/3hvp6ypR64Az/bAiA6ohzA472QXH9j3H8kACwzXhWMfYX+m6XuREJ4WqGDlg==","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"},{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEYCIQCGQZM1AIyo9lycAm9hTjkp6jLrNjiRPvozEoWp1uGVLgIhAIWC5br6qevQlhC97VowndApnm+pd+FGqJAa8PuR0kGc"}],"unpackedSize":613241},"main":"./dist/index.js","name":"@anyq/kafka","type":"module","types":"./dist/index.d.ts","author":"Shantanu Sharma","shasum":"186e7d558e69504fa97ec3cc7ed02c109f55723c","exports":{".":{"types":"./dist/index.d.ts","import":"./dist/index.js","require":"./dist/index.js"}},"license":"Apache-2.0","scripts":{"test":"bun test","build":"bun build ./src/index.ts --outdir ./dist --target bun && bun run build:types","clean":"rm -rf dist","typecheck":"tsc --noEmit","build:types":"tsc --emitDeclarationOnly --outDir dist"},"version":"0.5.0","_npmUser":{"name":"sns45","email":"notifyshantanu@gmail.com"},"homepage":"https://github.com/sns45/anyq#readme","keywords":["queue","kafka","message-queue","streaming","anyq"],"_integrity":"sha512-6RYjuIkhlHMk8qxSWKrsinGRgEvFZPVrKrWXBze8QcCzev3Ot+8qxT1Z4oVcgYzPoxQY+1mF20OdRT90TZ4Xkw==","repository":{"url":"https://github.com/sns45/anyq","type":"git","directory":"packages/kafka"},"_npmVersion":"10.8.3","description":"Apache Kafka adapter for anyq","directories":{},"maintainers":[{"name":"sns45","email":"notifyshantanu@gmail.com"}],"_nodeVersion":"26.3.0","dependencies":{"kafkajs":"^2.2.4","@anyq/core":"0.5.0"},"publishConfig":{"access":"public"},"_hasShrinkwrap":false,"devDependencies":{"@types/bun":"^1.1.0","typescript":"^5.0.0"},"peerDependencies":{"@anyq/core":"0.5.0"},"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/kafka_0.5.0_1788809125038_0.470812444783435"}}},"time":{"created":"2025-12-17T17:02:48.892Z","modified":"2026-09-07T19:25:25.423Z","0.1.0":"2025-12-17T17:02:49.143Z","0.1.1":"2025-12-17T17:09:43.053Z","0.3.0":"2026-06-20T13:26:19.447Z","0.3.1":"2026-06-21T15:10:03.185Z","0.3.2":"2026-07-24T16:24:30.642Z","0.4.0":"2026-08-27T19:56:14.626Z","0.5.0":"2026-09-07T19:25:25.245Z"},"bugs":{"url":"https://github.com/sns45/anyq/issues"},"author":"Shantanu Sharma","license":"Apache-2.0","homepage":"https://github.com/sns45/anyq#readme","keywords":["queue","kafka","message-queue","streaming","anyq"],"repository":{"url":"https://github.com/sns45/anyq","type":"git","directory":"packages/kafka"},"description":"Apache Kafka adapter for anyq","maintainers":[{"name":"sns45","email":"notifyshantanu@gmail.com"}],"readme":"# @anyq/kafka\n\nApache Kafka adapter for **anyq** - high-throughput distributed streaming.\n\n## Installation\n\n```bash\nnpm install @anyq/kafka @anyq/core kafkajs\n```\n\n## Usage\n\n```typescript\nimport { KafkaProducer, KafkaConsumer } from '@anyq/kafka';\n\n// Create producer\nconst producer = new KafkaProducer({\n  brokers: ['localhost:9092'],\n  topic: 'my-topic',\n  clientId: 'my-app'\n});\n\n// Create consumer\nconst consumer = new KafkaConsumer({\n  brokers: ['localhost:9092'],\n  topic: 'my-topic',\n  groupId: 'my-consumer-group',\n  clientId: 'my-app'\n});\n\nawait producer.connect();\nawait consumer.connect();\n\n// Subscribe to messages\nawait consumer.subscribe(async (message) => {\n  console.log('Received:', message.data);\n  console.log('Partition:', message.metadata.partition);\n  console.log('Offset:', message.metadata.offset);\n  await message.ack();\n});\n\n// Publish messages\nawait producer.publish({\n  event: 'click',\n  userId: 'user-123'\n});\n\n// Publish with key (for partitioning)\nawait producer.publish(\n  { event: 'purchase' },\n  { key: 'user-123' }  // Same key = same partition\n);\n\n// Cleanup\nawait consumer.disconnect();\nawait producer.disconnect();\n```\n\n## Configuration\n\n```typescript\ninterface KafkaConfig {\n  brokers: string[];         // Kafka broker addresses\n  topic: string;             // Topic name\n  clientId: string;          // Client identifier\n  groupId?: string;          // Consumer group (consumer only)\n  // SSL/SASL\n  ssl?: boolean;\n  sasl?: {\n    mechanism: 'plain' | 'scram-sha-256' | 'scram-sha-512';\n    username: string;\n    password: string;\n  };\n  // Producer options\n  acks?: -1 | 0 | 1;         // Acknowledgment level\n  compression?: 'gzip' | 'snappy' | 'lz4';\n  // Consumer options\n  fromBeginning?: boolean;   // Start from earliest offset\n  autoCommit?: boolean;      // Auto-commit offsets\n}\n```\n\n## Features\n\n- Consumer groups for load balancing\n- Partition-based ordering (by key)\n- Compression (gzip, snappy, lz4)\n- SSL/SASL authentication\n- Idempotent producer\n- Manual offset management\n\n## Retry strategies (0.3.0)\n\nThis adapter participates in the opt-in pluggable retry strategies from `@anyq/core`. Pass a strategy via `BaseQueueConfig.strategy` and `KafkaConsumer.applyStrategy()` takes over per-message error handling; omit it and the legacy catch behaviour (log + emit `error`; with `autoCommit` on, effectively skips the message) runs unchanged.\n\n| Capability | Support |\n|---|---|\n| `supportsNativeDelay` | **false** (this release) |\n| `park` | downgrades to in-process retry with a `warn` log, capped by `maxAttempts`. The tiered retry-topics implementation is a follow-up |\n| `deadLetterMessage` | default (`nack(false)` is a no-op for kafka; with `autoCommit` the offset advances past the message). Publish to a `<topic>.dlq` from your handler if you need a real DLQ |\n| Batch (`subscribeBatch`) | Kafka commits offsets per-batch, so handler failure is **all-or-nothing**: the strategy is applied to the first message as the batch representative |\n\n```typescript\nimport { logAndSkip } from '@anyq/core';\n\nconst consumer = new KafkaConsumer({\n  driver: 'kafka',\n  kafka: { brokers: ['localhost:9092'] },\n  topic: 'orders',\n  consumerGroup: { groupId: 'order-processors' },\n  strategy: logAndSkip(),\n});\n```\n\nSee `@anyq/core` for the full strategy catalogue.\n\n## License\n\nApache License 2.0. See [LICENSE](https://github.com/sns45/anyq/blob/main/LICENSE).\n","readmeFilename":"README.md"}