Loading skill
Install any skill in seconds. Free to start, no credit card required.
Get Started Free →Build event streaming applications using Azure Event Hubs SDK for JavaScript (@azure/event-hubs). Use when implementing high-throughput event ingestion, real-time analytics, IoT telemetry, or event-driven architectures with partitioned consumers.
| Test case | Without → With | Effect | Δ tokens | Δ turns |
|---|---|---|---|---|
| case-01 | ✗→✓ | ▲ Improved | 18% | 0% |
| case-02 | ✗→✓ | ▲ Improved | 25% | 0% |
| case-03 | ✗→✓ | ▲ Improved | 30% | 0% |
| case-07 | ✗→✓ | ▲ Improved | 60% | 0% |
| case-14 | ✗→✓ | ▲ Improved | 26% | 0% |
High-throughput event streaming and real-time data ingestion.
bashnpm install @azure/event-hubs @azure/identity
For checkpointing with consumer groups:
bashnpm install @azure/eventhubs-checkpointstore-blob @azure/storage-blob
bashEVENTHUB_NAMESPACE=<namespace>.servicebus.windows.net EVENTHUB_NAME=my-eventhub STORAGE_ACCOUNT_NAME=<storage-account> STORAGE_CONTAINER_NAME=checkpoints AZURE_TOKEN_CREDENTIALS=prod # Required only if DefaultAzureCredential is used in production
typescriptimport { EventHubProducerClient, EventHubConsumerClient } from "@azure/event-hubs"; import { DefaultAzureCredential, ManagedIdentityCredential } from "@azure/identity"; const fullyQualifiedNamespace = process.env.EVENTHUB_NAMESPACE!; const eventHubName = process.env.EVENTHUB_NAME!; // Local dev: DefaultAzureCredential. Production: set AZURE_TOKEN_CREDENTIALS=prod or AZURE_TOKEN_CREDENTIALS=<specific_credential> const credential = new DefaultAzureCredential({requiredEnvVars: ["AZURE_TOKEN_CREDENTIALS"]}); // Or use a specific credential directly in production: // See https://learn.microsoft.com/javascript/api/overview/azure/identity-readme?view=azure-node-latest#credential-classes // const credential = new ManagedIdentityCredential(); // Producer const producer = new EventHubProducerClient(fullyQualifiedNamespace, eventHubName, credential); // Consumer const consumer = new EventHubConsumerClient( "$Default", // Consumer group fullyQualifiedNamespace, eventHubName, credential );
typescriptconst producer = new EventHubProducerClient(namespace, eventHubName, credential); // Create batch and add events const batch = await producer.createBatch(); batch.tryAdd({ body: { temperature: 72.5, deviceId: "sensor-1" } }); batch.tryAdd({ body: { temperature: 68.2, deviceId: "sensor-2" } }); await producer.sendBatch(batch); await producer.close();
typescript// By partition ID const batch = await producer.createBatch({ partitionId: "0" }); // By partition key (consistent hashing) const batch = await producer.createBatch({ partitionKey: "device-123" });
typescriptconst consumer = new EventHubConsumerClient("$Default", namespace, eventHubName, credential); const subscription = consumer.subscribe({ processEvents: async (events, context) => { for (const event of events) { console.log(`Partition: ${context.partitionId}, Body: ${JSON.stringify(event.body)}`); } }, processError: async (err, context) => { console.error(`Error on partition ${context.partitionId}: ${err.message}`); }, }); // Stop after some time setTimeout(async () => { await subscription.close(); await consumer.close(); }, 60000);
typescriptimport { EventHubConsumerClient } from "@azure/event-hubs"; import { ContainerClient } from "@azure/storage-blob"; import { BlobCheckpointStore } from "@azure/eventhubs-checkpointstore-blob"; const containerClient = new ContainerClient( `https://${storageAccount}.blob.core.windows.net/${containerName}`, credential ); const checkpointStore = new BlobCheckpointStore(containerClient); const consumer = new EventHubConsumerClient( "$Default", namespace, eventHubName, credential, checkpointStore ); const subscription = consumer.subscribe({ processEvents: async (events, context) => { for (const event of events) { console.log(`Processing: ${JSON.stringify(event.body)}`); } // Checkpoint after processing batch if (events.length > 0) { await context.updateCheckpoint(events[events.length - 1]); } }, processError: async (err, context) => { console.error(`Error: ${err.message}`); }, });
typescriptconst subscription = consumer.subscribe({ processEvents: async (events, context) => { /* ... */ }, processError: async (err, context) => { /* ... */ }, }, { startPosition: { // Start from beginning "0": { offset: "@earliest" }, // Start from end (new events only) "1": { offset: "@latest" }, // Start from specific offset "2": { offset: "12345" }, // Start from specific time "3": { enqueuedOn: new Date("2024-01-01") }, }, });
typescript// Get hub info const hubProperties = await producer.getEventHubProperties(); console.log(`Partitions: ${hubProperties.partitionIds}`); // Get partition info const partitionProperties = await producer.getPartitionProperties("0"); console.log(`Last sequence: ${partitionProperties.lastEnqueuedSequenceNumber}`);
typescriptconst subscription = consumer.subscribe( { processEvents: async (events, context) => { /* ... */ }, processError: async (err, context) => { /* ... */ }, }, { maxBatchSize: 100, // Max events per batch maxWaitTimeInSeconds: 30, // Max wait for batch } );
typescriptimport { EventHubProducerClient, EventHubConsumerClient, EventData, ReceivedEventData, PartitionContext, Subscription, SubscriptionEventHandlers, CreateBatchOptions, EventPosition, } from "@azure/event-hubs"; import { BlobCheckpointStore } from "@azure/eventhubs-checkpointstore-blob";
typescript// Send with properties const batch = await producer.createBatch(); batch.tryAdd({ body: { data: "payload" }, properties: { eventType: "telemetry", deviceId: "sensor-1", }, contentType: "application/json", correlationId: "request-123", }); // Access in receiver consumer.subscribe({ processEvents: async (events, context) => { for (const event of events) { console.log(`Type: ${event.properties?.eventType}`); console.log(`Sequence: ${event.sequenceNumber}`); console.log(`Enqueued: ${event.enqueuedTimeUtc}`); console.log(`Offset: ${event.offset}`); } }, });
typescriptconsumer.subscribe({ processEvents: async (events, context) => { try { for (const event of events) { await processEvent(event); } await context.updateCheckpoint(events[events.length - 1]); } catch (error) { // Don't checkpoint on error - events will be reprocessed console.error("Processing failed:", error); } }, processError: async (err, context) => { if (err.name === "MessagingError") { // Transient error - SDK will retry console.warn("Transient error:", err.message); } else { // Fatal error console.error("Fatal error:", err); } }, });
createBatch() for efficient sendinglastEnqueuedSequenceNumber vs processed sequenceOther measured skills in the registry, with their headline benchmark lift.