{"_id":"@electreonwireless/native-aws-kcl-ts","_rev":"3-4017a936646d2ff31aaf055e6fb8beb3","name":"@electreonwireless/native-aws-kcl-ts","dist-tags":{"latest":"3.0.2"},"versions":{"3.0.1":{"name":"@electreonwireless/native-aws-kcl-ts","version":"3.0.1","keywords":["kinesis","aws","kcl","consumer","streaming"],"author":"","license":"UNLICENSED","_id":"@electreonwireless/native-aws-kcl-ts@3.0.1","maintainers":[{"name":"kosta-electreon","email":"kosta.k@electreon.com"}],"homepage":"https://github.com/electreonwireless/native-aws-kcl-ts#readme","bugs":{"url":"https://github.com/electreonwireless/native-aws-kcl-ts/issues"},"dist":{"shasum":"5fd524b4b6d6bda798a5da743c2433432e6cf1ed","tarball":"https://registry.npmjs.org/@electreonwireless/native-aws-kcl-ts/-/native-aws-kcl-ts-3.0.1.tgz","fileCount":247,"integrity":"sha512-7ThshPg7WV2ueEn6LtCBcbAS2/JdC11B+1wTbMH4WKBmkmeHSMk3ZrmIG07AwBcrxJ5xxhe9aJC9vGjflOHx2A==","signatures":[{"sig":"MEUCIEdkre61wpH7OFJSzitsY8pC42YFLF45iUfVY4liSCXmAiEA6Ls5HW19nNR7DXVyRxTaG4lkUaasBUxvak+sAfLXp64=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":1533301},"main":"dist/index.js","_from":"file:electreonwireless-native-aws-kcl-ts-3.0.1.tgz","types":"dist/index.d.ts","engines":{"node":">=24.0.0"},"private":false,"scripts":{"test":"jest --passWithNoTests","build":"tsc -p tsconfig.build.json","release":"pnpm run build && changeset publish","prebuild":"rimraf dist","version-packages":"changeset version"},"_npmUser":{"name":"kosta-electreon","email":"kosta.k@electreon.com"},"_resolved":"/tmp/a30f86130825bed0e7b7860099917089/electreonwireless-native-aws-kcl-ts-3.0.1.tgz","_integrity":"sha512-7ThshPg7WV2ueEn6LtCBcbAS2/JdC11B+1wTbMH4WKBmkmeHSMk3ZrmIG07AwBcrxJ5xxhe9aJC9vGjflOHx2A==","repository":{"url":"git+https://github.com/electreonwireless/native-aws-kcl-ts.git","type":"git"},"_npmVersion":"11.9.0","description":"Native TypeScript implementation of AWS Kinesis Client Library (KCL) with async/await support","directories":{},"_nodeVersion":"24.14.0","dependencies":{"uuid":"^10.0.0","@aws-sdk/client-kinesis":"^3.700.0","@aws-sdk/client-dynamodb":"^3.700.0","@aws-sdk/client-cloudwatch":"^3.700.0","@aws-sdk/credential-providers":"^3.700.0"},"publishConfig":{"access":"public","provenance":false},"_hasShrinkwrap":false,"devDependencies":{"jest":"^30.3.0","rimraf":"^6.1.3","ts-jest":"^29.4.6","typescript":"^4.9.4","@types/jest":"28.1.8","@types/node":"^24.0.0","@types/uuid":"^10.0.0","@changesets/cli":"^2.30.0"},"peerDependencies":{"pg":"^8.0.0","mysql2":"^3.0.0","ioredis":"^5.0.0","mongodb":"^6.0.0","@aws-sdk/client-s3":"^3.700.0"},"peerDependenciesMeta":{"pg":{"optional":true},"mysql2":{"optional":true},"ioredis":{"optional":true},"mongodb":{"optional":true},"@aws-sdk/client-s3":{"optional":true}},"_npmOperationalInternal":{"tmp":"tmp/native-aws-kcl-ts_3.0.1_1777536557519_0.029571369314510987","host":"s3://npm-registry-packages-npm-production"}},"3.0.2":{"name":"@electreonwireless/native-aws-kcl-ts","version":"3.0.2","description":"Native TypeScript implementation of AWS Kinesis Client Library (KCL) with async/await support","author":"","private":false,"publishConfig":{"access":"public","provenance":false},"license":"UNLICENSED","main":"dist/index.js","types":"dist/index.d.ts","repository":{"type":"git","url":"git+https://github.com/electreonwireless/native-aws-kcl-ts.git"},"keywords":["kinesis","aws","kcl","consumer","streaming"],"dependencies":{"@aws-sdk/client-cloudwatch":"^3.700.0","@aws-sdk/client-dynamodb":"^3.700.0","@aws-sdk/client-kinesis":"^3.700.0","@aws-sdk/credential-providers":"^3.700.0","uuid":"^10.0.0"},"peerDependencies":{"@aws-sdk/client-s3":"^3.700.0","ioredis":"^5.0.0","mongodb":"^6.0.0","mysql2":"^3.0.0","pg":"^8.0.0"},"peerDependenciesMeta":{"@aws-sdk/client-s3":{"optional":true},"ioredis":{"optional":true},"pg":{"optional":true},"mysql2":{"optional":true},"mongodb":{"optional":true}},"devDependencies":{"@changesets/cli":"^2.30.0","@types/uuid":"^10.0.0","jest":"^30.3.0","rimraf":"^6.1.3","ts-jest":"^29.4.6","typescript":"^4.9.4","@types/jest":"28.1.8","@types/node":"^24.0.0"},"engines":{"node":">=24.0.0"},"scripts":{"prebuild":"rimraf dist","build":"tsc -p tsconfig.build.json","test":"jest --passWithNoTests","release":"pnpm run build && changeset publish","version-packages":"changeset version"},"_id":"@electreonwireless/native-aws-kcl-ts@3.0.2","bugs":{"url":"https://github.com/electreonwireless/native-aws-kcl-ts/issues"},"homepage":"https://github.com/electreonwireless/native-aws-kcl-ts#readme","_integrity":"sha512-b6zvUkLKgSuUcc9Cwr2C4avudPqp6DC5kQgMDDwCjMThRLamcO/hj27q2Ikum9b4WRX2Ci2y2ZYau5MSgut+iA==","_resolved":"/tmp/136c4e14627603760f2a4f9dcc2ab5d6/electreonwireless-native-aws-kcl-ts-3.0.2.tgz","_from":"file:electreonwireless-native-aws-kcl-ts-3.0.2.tgz","_nodeVersion":"24.14.1","_npmVersion":"11.13.0","dist":{"integrity":"sha512-b6zvUkLKgSuUcc9Cwr2C4avudPqp6DC5kQgMDDwCjMThRLamcO/hj27q2Ikum9b4WRX2Ci2y2ZYau5MSgut+iA==","shasum":"e88c62d45a15f4cbdf3f05894a48a85d57d4e1f1","tarball":"https://registry.npmjs.org/@electreonwireless/native-aws-kcl-ts/-/native-aws-kcl-ts-3.0.2.tgz","fileCount":247,"unpackedSize":1533586,"signatures":[{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEYCIQDHHH6jFDnr7XqP65JiPDnlXCCq6hmWcDxRbl2XGuLUKAIhAMeODCRouOLIx/QpVA+HzFfhvgEGfu0FCcNwcM75QlkO"}]},"_npmUser":{"name":"GitHub Actions","email":"npm-oidc-no-reply@github.com","trustedPublisher":{"id":"github","oidcConfigId":"oidc:99179531-3139-4391-b49b-b1206cd1ba7f"}},"directories":{},"maintainers":[{"name":"kosta-electreon","email":"kosta.k@electreon.com"},{"name":"electreon-npm","email":"omri.s@electreon.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/native-aws-kcl-ts_3.0.2_1777546113427_0.3604081600088993"},"_hasShrinkwrap":false}},"time":{"created":"2026-04-30T08:09:17.427Z","modified":"2026-04-30T10:48:34.096Z","3.0.1":"2026-04-30T08:09:17.735Z","3.0.2":"2026-04-30T10:48:33.623Z"},"bugs":{"url":"https://github.com/electreonwireless/native-aws-kcl-ts/issues"},"license":"UNLICENSED","homepage":"https://github.com/electreonwireless/native-aws-kcl-ts#readme","keywords":["kinesis","aws","kcl","consumer","streaming"],"repository":{"type":"git","url":"git+https://github.com/electreonwireless/native-aws-kcl-ts.git"},"description":"Native TypeScript implementation of AWS Kinesis Client Library (KCL) with async/await support","maintainers":[{"name":"kosta-electreon","email":"kosta.k@electreon.com"},{"name":"electreon-npm","email":"omri.s@electreon.com"}],"readme":"# @electreonwireless/native-aws-kcl-ts\n\nA native TypeScript implementation of the AWS Kinesis Client Library (KCL) for consuming Kinesis streams. This library provides the same functionality as the Java-based KCL but is written entirely in TypeScript with an `async/await` concurrency model.\n\n## Features\n\n- **Shard Discovery**: Automatically discovers and tracks Kinesis shards\n- **Lease Management**: Coordinates shard-to-worker assignment using pluggable storage (DynamoDB, Redis, S3, PostgreSQL, MySQL, MongoDB)\n- **Lease Renewal**: Maintains heartbeats to prevent lease expiration\n- **Failover & Rebalancing**: Automatically handles worker failures and redistributes shards\n- **Checkpoint Coordination**: Persists processing progress with pluggable storage backends\n- **Graceful Shutdown**: Handles shutdown signals and shard end events\n- **Enhanced Fan-Out (EFO)**: Optional dedicated throughput mode with 2MB/s per consumer per shard\n- **Exponential Backoff with Jitter**: Handles throttling gracefully following AWS best practices\n- **Pluggable Logging**: Configure your own logger factory or use the default console logger\n\n## Installation\n\n### Base Installation (DynamoDB - Default)\n\n```bash\npnpm add @electreonwireless/native-aws-kcl-ts\n```\n\n### With Optional Persistence Backends\n\nInstall the base package plus the driver for your chosen persistence backend:\n\n```bash\n# S3 persistence\npnpm add @electreonwireless/native-aws-kcl-ts @aws-sdk/client-s3\n\n# Redis persistence\npnpm add @electreonwireless/native-aws-kcl-ts ioredis\n\n# PostgreSQL persistence\npnpm add @electreonwireless/native-aws-kcl-ts pg\n\n# MySQL persistence\npnpm add @electreonwireless/native-aws-kcl-ts mysql2\n\n# MongoDB persistence\npnpm add @electreonwireless/native-aws-kcl-ts mongodb\n```\n\n> **Note**: The persistence drivers (`ioredis`, `pg`, `mysql2`, `mongodb`, `@aws-sdk/client-s3`) are optional peer dependencies. Only install the one you need. If you don't install a driver and try to use that persistence type, you'll get a clear error message telling you which package to install.\n\n## Quick Start\n\n### Basic Usage\n\n```typescript\nimport { \n  Scheduler, \n  ConfigBuilder,\n  BaseRecordProcessor,\n  processorFactory,\n} from '@electreonwireless/native-aws-kcl-ts';\n\n// Create a custom record processor\nclass MyProcessor extends BaseRecordProcessor {\n  protected async processRecord(record: KinesisClientRecord): Promise<void> {\n    const data = Buffer.from(record.data).toString('utf-8');\n    console.log('Processing record:', data);\n    // Your processing logic here\n  }\n}\n\n// Build configuration\nconst config = new ConfigBuilder()\n  .withAwsRegion('us-east-1')\n  .withStreamName('my-stream')\n  .withApplicationName('my-app')\n  .withProcessorFactory(processorFactory(MyProcessor))\n  .withInitialPosition('TRIM_HORIZON')\n  .build();\n\n// Create and run the scheduler\nconst scheduler = new Scheduler(config);\nawait scheduler.run();\n\n// Graceful shutdown\nprocess.on('SIGTERM', async () => {\n  await scheduler.shutdown();\n});\n```\n\n### Configuring a Custom Logger\n\nThe library uses a simple console logger by default. You can configure a custom logger factory:\n\n```typescript\nimport { configureNativeKclLibrary } from '@electreonwireless/native-aws-kcl-ts';\nimport winston from 'winston';\n\nconst myWinston = winston.createLogger({\n  level: 'info',\n  format: winston.format.json(),\n  transports: [new winston.transports.Console()],\n});\n\n// Configure the library once at startup\nconfigureNativeKclLibrary({\n  loggerFactory: (name: string) => myWinston.child({ context: name })\n});\n\n// Now all classes in the library will use your Winston logger\nconst scheduler = new Scheduler(config);\n```\n\n## Core Concepts\n\n### Scheduler\n\nThe `Scheduler` (also known as Worker in KCL) is the main entry point. It coordinates all components:\n- LeaseCoordinator: Manages shard leases for distributed processing\n- ShardDetector: Discovers shards in the stream\n- ShardSyncer: Synchronizes shards with leases\n- ShardConsumers: Process records from individual shards\n\n### Record Processors\n\nRecord processors implement the `ShardRecordProcessor` interface to handle records from a shard. You can either:\n- Extend `BaseRecordProcessor` for a simplified interface\n- Implement `ShardRecordProcessor` directly for full control\n\n### Lease Management\n\nLeases coordinate which worker processes which shard. The library supports multiple storage backends:\n- **DynamoDB** (default)\n- **Redis**\n- **S3**\n- **PostgreSQL**\n- **MySQL**\n- **MongoDB**\n\n## Enhanced Fan-Out (EFO) Mode\n\nEnhanced Fan-Out provides dedicated throughput of 2 MB/s per consumer per shard, compared to the shared 2 MB/s per shard with polling mode.\n\n### Comparison of Retrieval Modes\n\n| Feature | Polling (DEFAULT) | Enhanced Fan-Out (FANOUT) |\n|---------|-------------------|---------------------------|\n| Throughput | Shared 2 MB/s per shard | Dedicated 2 MB/s per consumer per shard |\n| Latency | ~200ms+ | ~70ms |\n| API | GetRecords | SubscribeToShard (HTTP/2 push) |\n| Consumers per shard | Limited by shared throughput | Up to 20 (or 50 with EFO Advantage) |\n| Cost | Standard Kinesis pricing | Additional per-consumer charges |\n\n### Enabling Enhanced Fan-Out\n\n```typescript\nconst config = new ConfigBuilder()\n  .withAwsRegion('us-east-1')\n  .withStreamName('my-stream')\n  .withApplicationName('my-app')\n  .withProcessorFactory(processorFactory(MyProcessor))\n  .withRetrievalMode('FANOUT') // Enable EFO\n  .build();\n```\n\n## Pluggable Persistence Layer\n\nBy default, the library uses **DynamoDB** for storing lease and checkpoint information. However, the persistence layer is abstracted behind the `LeaseRefresher` interface, allowing you to implement custom storage backends.\n\n### Supported Storage Backends\n\n| Backend | Status | Config Value | Extra Package |\n|---------|--------|--------------|---------------|\n| DynamoDB | ✅ Default | `dynamodb` | None (built-in) |\n| S3 | ✅ Supported | `s3` | `@aws-sdk/client-s3` |\n| Redis | ✅ Supported | `redis` | `ioredis` |\n| PostgreSQL | ✅ Supported | `postgresql` | `pg` |\n| MySQL | ✅ Supported | `mysql` | `mysql2` |\n| MongoDB | ✅ Supported | `mongodb` | `mongodb` |\n| Custom | ✅ Supported | N/A | User-provided |\n\n### Using a Custom Persistence Backend\n\n```typescript\nimport { \n  LeaseRefresher, \n  Lease,\n  KinesisConsumerConfig,\n  Scheduler,\n} from '@electreonwireless/native-aws-kcl-ts';\n\n// Implement the LeaseRefresher interface\nclass MyCustomLeaseRefresher implements LeaseRefresher {\n  // ... implement all required methods\n}\n\n// Use in configuration\nconst config: KinesisConsumerConfig = {\n  aws: { region: 'us-east-1' },\n  applicationName: 'my-app',\n  streamName: 'my-stream',\n  recordProcessorFactory: myFactory,\n  persistence: {\n    leaseRefresher: new MyCustomLeaseRefresher(),\n  },\n};\n\nconst scheduler = new Scheduler(config);\nawait scheduler.run();\n```\n\n## Metrics\n\nThe library supports publishing application metrics to AWS CloudWatch (or a custom monitoring backend). By default, metrics are **disabled** (level: `NONE`).\n\n### Metrics Levels\n\n| Level | Description |\n|-------|-------------|\n| `NONE` | No metrics are reported (default) |\n| `SUMMARY` | Aggregated metrics: RecordsProcessed, MillisBehindLatest, LeasesHeld |\n| `DETAILED` | All SUMMARY metrics plus per-shard and per-operation metrics |\n\n### Enabling Metrics\n\n```typescript\nconst config = new ConfigBuilder()\n  .withAwsRegion('us-east-1')\n  .withStreamName('my-stream')\n  .withApplicationName('my-app')\n  .withProcessorFactory(processorFactory(MyProcessor))\n  .withMetrics('DETAILED', 'MyApp/KCL')\n  .build();\n```\n\n### Custom Metrics Backend\n\nYou can implement your own metrics backend (e.g., Prometheus, Datadog) by implementing the `IMetricsFactory` interface:\n\n```typescript\nimport { \n  IMetricsFactory, \n  IMetricsScope,\n  MetricsLevel,\n} from '@electreonwireless/native-aws-kcl-ts';\n\nclass MyPrometheusMetricsFactory implements IMetricsFactory {\n  createScope(operation: string): IMetricsScope {\n    return new MyPrometheusScope(operation);\n  }\n  \n  getMetricsLevel(): MetricsLevel {\n    return MetricsLevel.DETAILED;\n  }\n  \n  isEnabled(): boolean {\n    return true;\n  }\n  \n  isDetailedMetricsEnabled(): boolean {\n    return true;\n  }\n  \n  async shutdown(): Promise<void> {\n    // Flush metrics\n  }\n}\n\n// Use the custom factory\nconst config: KinesisConsumerConfig = {\n  // ... other config\n  metrics: {\n    level: 'DETAILED',\n    customFactory: new MyPrometheusMetricsFactory(),\n  },\n};\n```\n\n## API Reference\n\n### Scheduler\n\nThe main entry point for the KCL.\n\n```typescript\nclass Scheduler {\n  constructor(config: KinesisConsumerConfig);\n  async run(): Promise<void>;\n  async shutdown(): Promise<void>;\n  getState(): SchedulerState;\n  getActiveConsumerCount(): number;\n}\n```\n\n### ConfigBuilder\n\nA fluent builder for creating `KinesisConsumerConfig`:\n\n```typescript\nconst config = new ConfigBuilder()\n  .withAwsRegion('us-east-1')\n  .withStreamName('my-stream')\n  .withApplicationName('my-app')\n  .withProcessorFactory(processorFactory(MyProcessor))\n  .withInitialPosition('TRIM_HORIZON')\n  .withRetrievalMode('DEFAULT')\n  .withMetrics('SUMMARY', 'MyApp/KCL')\n  .build();\n```\n\n### BaseRecordProcessor\n\nA base class that simplifies implementing record processors:\n\n```typescript\nclass MyProcessor extends BaseRecordProcessor {\n  protected async processRecord(record: KinesisClientRecord): Promise<void> {\n    // Process individual records\n  }\n  \n  async initialize(input: InitializationInput): Promise<void> {\n    // Optional: Custom initialization\n    await super.initialize(input);\n  }\n  \n  async leaseLost(input: LeaseLostInput): Promise<void> {\n    // Optional: Handle lease loss\n  }\n  \n  async shardEnded(input: ShardEndedInput): Promise<void> {\n    // Optional: Handle shard end\n    await super.shardEnded(input); // Default checkpoints at SHARD_END\n  }\n}\n```\n\n## License\n\nUNLICENSED - Proprietary\n","readmeFilename":"README.md"}