Imported from paulpas/agent-skill-router (
skills/cncf/cloudevents/SKILL.md). Install upstream withnpx skills add paulpas/agent-skill-router --skill cloudevents. Copyright stays with the author (MIT).
CloudEvents in Cloud-Native Engineering
Category: eventing
Status: Active
Stars: 3,300
Last Updated: 2026-04-22
Primary Language: Multiple (spec-first)
Documentation: https://cloudevents.io/
Purpose and Use Cases
CloudEvents is a CNCF incubating project that provides a standardized way to define event data in a common format, enabling interoperability across different event systems and services.
What Problem Does It Solve?
Event vendor lock-in and inconsistent event formats across platforms. Before CloudEvents, each messaging system (Kafka, RabbitMQ, AWS SNS, etc.) used its own event structure, making it difficult to build portable event-driven applications.
When to Use This Project
Use CloudEvents when you need to:
- Build event-driven architectures that span multiple platforms
- Integrate events across different messaging systems
- Create portable event producers and consumers
- Implement standardized event schemas for audit and compliance
Key Use Cases
- Event Routing: Forward events between different event buses (Kafka → PubSub → SQS)
- Cross-Platform Integration: Connect SaaS platforms (GitHub, Slack) with internal systems
- Serverless Event Processing: AWS Lambda, Azure Functions, Cloud Functions all support CloudEvents
- Observability: Standardized event format for tracing and debugging distributed systems
- IoT Data Ingestion: Standardized format for device events across cloud providers
Architecture Design Patterns
Event Structure
CloudEvents defines a standard structure with required and optional attributes:
{
"specversion": "1.0",
"type": "com.example.someevent",
"source": "/mycontext/source",
"id": "A234-1234-1234",
"time": "2018-04-05T17:31:00Z",
"datacontenttype": "application/json",
"data": {
"appinfoA": "abc",
"appinfoB": "xyz"
}
}
Required Attributes:
specversion: CloudEvents specification versiontype: Event type (backwardslash-separated namespacing)source: Event origin (URI or URN)id: Unique identifier within the scope of source
Optional Attributes:
subject: Subject of the eventtime: Timestamp of when the event happeneddatacontenttype: MIME type of the data payloaddataschema: URI referencing the schema of datadata: Event payload (binary or JSON)
Event Consumption Patterns
Direct Invocation
// HTTP POST with CloudEvents in body
POST /events HTTP/1.1
Content-Type: application/cloudevents+json
{
"specversion": "1.0",
"type": "com.example.order.created",
"source": "https://orders.example.com",
"id": "order-123",
"data": { "orderId": "123", "customerId": "456" }
}
Binary Content Mode
// HTTP POST with CloudEvents as headers
POST /events HTTP/1.1
Ce-Specversion: 1.0
Ce-Type: com.example.order.created
Ce-Source: https://orders.example.com
Ce-Id: order-123
Content-Type: application/json
{ "orderId": "123", "customerId": "456" }
Structured Content Mode
// HTTP POST with CloudEvents as envelope
POST /events HTTP/1.1
Content-Type: application/cloudevents+json
{
"specversion": "1.0",
"type": "com.example.order.created",
"source": "https://orders.example.com",
"id": "order-123",
"datacontenttype": "application/json",
"data": { "orderId": "123" }
}
Event Routing and Transformation
// Using Knative Eventing with CloudEvents
// Event source → Broker → Trigger → Service
// Each step passes CloudEvents through the system
// Trigger example with filtering
apiVersion: eventing.knative.dev/v1
kind: Trigger
metadata:
name: my-trigger
spec:
broker: default
filter:
attributes:
type: com.example.order.created
source: https://orders.example.com
subscriber:
ref:
apiVersion: serving.knative.dev/v1
kind: Service
name: order-processor
Data Encoding
CloudEvents supports multiple data content types:
JSON Data:
{
"specversion": "1.0",
"type": "com.example.order.created",
"source": "https://orders.example.com",
"id": "order-123",
"datacontenttype": "application/json",
"data": {
"orderId": "123",
"items": [{ "productId": "A", "quantity": 2 }]
}
}
Binary Data (Base64 encoded):
{
"specversion": "1.0",
"type": "com.example.file.uploaded",
"source": "https://storage.example.com",
"id": "file-456",
"datacontenttype": "application/octet-stream",
"data_base64": "SGVsbG8gV29ybGQ=" // "Hello World" base64
}
Text Data:
{
"specversion": "1.0",
"type": "com.example.notification.sent",
"source": "https://notify.example.com",
"id": "notif-789",
"datacontenttype": "text/plain",
"data": "Your order has been shipped"
}
Integration Approaches
Event Mesh Integration
CloudEvents work seamlessly with event mesh architectures:
// Kafka producer with CloudEvents
const producer = new Kafka.Producer({
'bootstrap.servers': 'localhost:9092'
});
const event = {
specversion: '1.0',
type: 'com.example.transaction.processed',
source: 'https://payments.example.com',
id: 'txn-' + Date.now(),
time: new Date().toISOString(),
datacontenttype: 'application/json',
data: { transactionId: 'tx-123', amount: 99.99 }
};
producer.produce({
topic: 'transactions',
value: Buffer.from(JSON.stringify(event))
});
Serverless Platform Integration
AWS Lambda with CloudEvents:
// Lambda function receiving CloudEvents
exports.handler = async (event) => {
// Event may be CloudEvents envelope or raw data
const cloudEvent = event.hasOwnProperty('specversion')
? event
: parseCloudEventFromHeaders(event.headers);
console.log(`Processing event ${cloudEvent.id} of type ${cloudEvent.type}`);
// Process based on event type
switch (cloudEvent.type) {
case 'com.example.order.created':
return processOrderCreated(cloudEvent.data);
case 'com.example.order.cancelled':
return processOrderCancelled(cloudEvent.data);
}
};
Azure Functions with CloudEvents:
// Azure Function with Event Grid trigger (CloudEvents format)
module.exports = async function (context, event) {
context.log(`Received CloudEvent: ${event.id}`);
context.log(`Type: ${event.type}`);
context.log(`Source: ${event.source}`);
// Process CloudEvent
processCloudEvent(event);
};
Event Bus Integration
CloudEvent → NATS → CloudEvent:
// NATS JetStream with CloudEvents
const nc = await connect({ servers: 'nats://localhost:4222' });
const js = nc.jetstream();
// Publish CloudEvent to NATS
await js.publish("orders.*", nats.jsonEncoded(cloudEvent));
// Subscribe to CloudEvents with filtering
const sub = await js.subscribe("orders.*", { durable: "processor" });
for await (const msg of sub) {
const event = msg.json();
console.log(`Received: ${event.type} from ${event.source}`);
}
Kubernetes Native Integration
Knative Eventing with CloudEvents:
# Knative Service consuming CloudEvents
apiVersion: serving.knative.dev/v1
kind: Service
metadata:
name: event-processor
spec:
template:
spec:
containers:
- image: event-processor
env:
- name: CE_TYPES
value: "com.example.order.created,com.example.order.shipped"
Broker/Trigger Pattern:
# Event broker receiving CloudEvents
apiVersion: eventing.knative.dev/v1
kind: Broker
metadata:
name: default
namespace: default
# Trigger filtering CloudEvents by type
apiVersion: eventing.knative.dev/v1
kind: Trigger
metadata:
name: order-processor
spec:
broker: default
filter:
attributes:
type: com.example.order.created
subscriber:
ref:
apiVersion: serving.knative.dev/v1
kind: Service
name: order-service
Common Pitfalls and How to Avoid Them
1. Missing Required Attributes
Pitfall: Omitting required CloudEvents attributes causes validation failures.
// ❌ Incorrect - missing required attributes
const badEvent = { type: 'test', data: {} };
// ✅ Correct - all required attributes present
const goodEvent = {
specversion: '1.0',
type: 'com.example.test',
source: 'https://example.com',
id: 'unique-id',
data: {}
};
2. Event ID Collisions
Pitfall: Reusing event IDs within the same source scope causes downstream deduplication issues.
// ❌ Incorrect - ID not unique
const badEvent = {
specversion: '1.0',
type: 'com.example.transaction',
source: 'https://payments.example.com',
id: 'txn-123', // Not unique!
data: {}
};
// ✅ Correct - UUID-based unique ID
const goodEvent = {
specversion: '1.0',
type: 'com.example.transaction',
source: 'https://payments.example.com',
id: 'txn-' + crypto.randomUUID(),
data: {}
};
3. Time Format Inconsistencies
Pitfall: Using non-RFC 3339 time formats causes parsing issues.
// ❌ Incorrect - non-standard time format
const badEvent = {
specversion: '1.0',
time: '2024-01-15 10:30:00', // Invalid format
data: {}
};
// ✅ Correct - RFC 3339 format
const goodEvent = {
specversion: '1.0',
time: '2024-01-15T10:30:00Z', // RFC 3339
data: {}
};
4. Type Naming Conventions
Pitfall: Inconsistent type naming makes event discovery and filtering difficult.
// ❌ Incorrect - inconsistent naming
const badEvent = {
type: 'orderCreated', // camelCase, no namespace
data: {}
};
// ✅ Correct - standardized naming
const goodEvent = {
type: 'com.example.order.created', // backwardslash, namespace-based
data: {}
};
5. Data Schema Mismatch
Pitfall: Data content type doesn't match actual data format.
// ❌ Incorrect - content type mismatch
const badEvent = {
datacontenttype: 'application/json',
data: '<xml><order>...</order></xml>' // Not JSON!
};
// ✅ Correct - matching content type
const goodEvent = {
datacontenttype: 'application/xml',
data: '<xml><order>...</order></xml>'
};
// Or properly converted to JSON
const goodEventJSON = {
datacontenttype: 'application/json',
data: { order: { id: '123', status: 'created' } }
};
6. Overly Large Data Payloads
Pitfall: Including large payloads in CloudEvents can cause performance issues.
// ❌ Incorrect - large embedded data
const badEvent = {
type: 'com.example.file.uploaded',
data: {
filename: 'large-file.zip',
content: base64LargeFile // 10MB embedded
}
};
// ✅ Correct - reference to data
const goodEvent = {
type: 'com.example.file.uploaded',
data: {
filename: 'large-file.zip',
url: 'https://storage.example.com/files/abc123'
}
};
Coding Practices
Event Producer Patterns
1. Event Factory Pattern
// Event factory for consistent CloudEvent creation
class CloudEventFactory {
constructor(source, defaultTypePrefix = '') {
this.source = source;
this.defaultTypePrefix = defaultTypePrefix;
}
create(type, data, attributes = {}) {
return {
specversion: '1.0',
type: this.defaultTypePrefix + type,
source: this.source,
id: this.generateId(),
time: new Date().toISOString(),
...attributes,
data: data
};
}
generateId() {
return crypto.randomUUID();
}
}
// Usage
const factory = new CloudEventFactory('https://orders.example.com', 'com.example.');
const event = factory.create('order.created', { orderId: '123' });
2. Decorator Pattern for Enrichment
// Decorate events with common attributes
function enrichEvent(event, context) {
return {
...event,
time: context.timestamp || new Date().toISOString(),
source: context.source || event.source,
subject: context.subject || event.subject,
extensions: {
...event.extensions,
correlationId: context.correlationId,
traceId: context.traceId
}
};
}
// Usage
const baseEvent = factory.create('order.shipped', { orderId: '123' });
const enrichedEvent = enrichEvent(baseEvent, {
correlationId: 'corr-abc',
traceId: 'trace-xyz'
});
3. Validation Pattern
// Validate CloudEvents
function validateCloudEvent(event) {
const errors = [];
// Check required attributes
const required = ['specversion', 'type', 'source', 'id'];
for (const attr of required) {
if (!event[attr]) {
errors.push(`Missing required attribute: ${attr}`);
}
}
// Validate specversion
if (event.specversion && event.specversion !== '1.0') {
errors.push(`Unsupported specversion: ${event.specversion}`);
}
// Validate type format
if (event.type && !event.type.includes('.')) {
errors.push('Type should contain at least one dot (backwardslash)');
}
// Validate time format (RFC 3339)
if (event.time && !/^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}(\.\d+)?(Z|[+-]\d{2}:\d{2})$/.test(event.time)) {
errors.push('Time must be in RFC 3339 format');
}
return errors.length === 0 ? null : errors;
}
// Usage
const errors = validateCloudEvent(event);
if (errors) {
throw new Error(`Invalid CloudEvent: ${errors.join(', ')}`);
}
Event Consumer Patterns
1. Event Router Pattern
// Event router for type-based handling
class CloudEventRouter {
constructor() {
this.handlers = new Map();
}
register(type, handler) {
const types = Array.isArray(type) ? type : [type];
for (const t of types) {
this.handlers.set(t, handler);
}
return this;
}
async route(event) {
const handler = this.handlers.get(event.type);
if (!handler) {
throw new Error(`No handler registered for event type: ${event.type}`);
}
return await handler(event);
}
}
// Usage
const router = new CloudEventRouter()
.register('com.example.order.created', processOrderCreated)
.register('com.example.order.shipped', processOrderShipped)
.register(['com.example.order.delivered', 'com.example.order.completed'], processOrderFinalized);
// Process incoming CloudEvent
await router.route(cloudEvent);
2. Filter Pattern
// Filter CloudEvents by attributes
function filterCloudEvent(event, filters) {
for (const [key, value] of Object.entries(filters)) {
// Check attribute
if (event[key] && event[key] !== value) {
return false;
}
// Check data path
if (key.includes('.') && event.data) {
const path = key.split('.');
const actual = path.reduce((obj, p) => obj?.[p], event.data);
if (actual !== value) {
return false;
}
}
}
return true;
}
// Usage
const filters = {
type: 'com.example.order.created',
'data.status': 'pending'
};
if (filterCloudEvent(event, filters)) {
processOrder(event.data);
}
3. Batch Processing Pattern
// Batch CloudEvents for efficiency
async function batchProcessEvents(events, batchSize = 100) {
const batches = [];
for (let i = 0; i < events.length; i += batchSize) {
batches.push(events.slice(i, i + batchSize));
}
for (const batch of batches) {
try {
await processBatch(batch);
} catch (error) {
// Handle batch error (retry, dead letter, etc.)
console.error(`Batch failed: ${error.message}`);
}
}
}
async function processBatch(batch) {
// Process batch of CloudEvents
const promises = batch.map(event => processSingleEvent(event));
return await Promise.all(promises);
}
Testing Patterns
1. Unit Test with Mock Events
// Test CloudEvent producer
describe('OrderEventProducer', () => {
it('should create valid CloudEvent', () => {
const factory = new CloudEventFactory('https://orders.example.com');
const event = factory.create('order.created', { orderId: '123' });
// Validate structure
expect(event.specversion).toBe('1.0');
expect(event.type).toBe('com.example.order.created');
expect(event.source).toBe('https://orders.example.com');
expect(event.id).toBeDefined();
expect(event.time).toBeDefined();
expect(event.data.orderId).toBe('123');
// Validate format
expect(validateCloudEvent(event)).toBeNull();
});
});
2. Integration Test with Event Bus
// Test end-to-end CloudEvent flow
describe('CloudEvent Integration', () => {
let eventBus;
let consumerSpy;
beforeEach(() => {
eventBus = new TestEventBus();
consumerSpy = jest.fn();
});
it('should deliver CloudEvent through event bus', async () => {
// Subscribe consumer
eventBus.subscribe('orders', consumerSpy);
// Publish CloudEvent
const event = factory.create('order.created', { orderId: '123' });
await eventBus.publish('orders', event);
// Verify delivery
await wait(100); // Wait for async delivery
expect(consumerSpy).toHaveBeenCalledWith(event);
});
});
Fundamentals
CloudEvents Specification Versions
Current Version: 1.0 (Stable)
CloudEvents 1.0 is the current stable specification, approved by the CNCF Technical Oversight Committee. It defines:
- Standard event structure with required and optional attributes
- Data encoding formats (JSON, binary, text)
- Transport bindings (HTTP, Kafka, AMQP, etc.)
- Attribute naming conventions
Attribute Types
Context Attributes:
specversion: String, requiredtype: String, required (backwardslash-separated namespacing)source: String, required (URI or URN)id: String, required (unique within source scope)subject: String, optionaltime: String, optional (RFC 3339 timestamp)
Data Attributes:
datacontenttype: String, optional (MIME type)dataschema: String, optional (URI to schema)data: Any, optional (event payload)
Extension Attributes:
- Custom attributes can be added for implementation-specific needs
- Must follow naming conventions (alphanumeric, hyphens, dots)
Protocol Bindings
HTTP Binding:
- CloudEvents can be sent as HTTP POST requests
- Supports both binary and structured content modes
- Common Content-Types:
application/cloudevents+json,application/json
Kafka Binding:
- Event attributes mapped to Kafka record headers
- Event data is the Kafka record value
- Type mapping: Kafka topic → CloudEvents type
AMQP Binding:
- Attributes mapped to AMQP 1.0 properties
- Event data in message body
- Type mapping: AMQP subject → CloudEvents type
Data Formats
JSON Data:
{
"specversion": "1.0",
"type": "com.example.order.created",
"source": "https://orders.example.com",
"id": "unique-id",
"datacontenttype": "application/json",
"data": { "orderId": "123" }
}
Binary Data (base64):
{
"specversion": "1.0",
"datacontenttype": "application/octet-stream",
"data_base64": "SGVsbG8="
}
Text Data:
{
"specversion": "1.0",
"datacontenttype": "text/plain",
"data": "Hello World"
}
Scaling and Deployment Patterns
Event Bus Architecture
Horizontal Scaling
// Scale event consumers horizontally
// Each consumer instance processes a subset of events
// Using Kafka consumer groups
const consumer = kafka.consumer({ groupId: 'order-processor' });
await consumer.connect();
await consumer.subscribe({ topic: 'orders' });
await consumer.run({
eachMessage: async ({ topic, partition, message }) => {
const event = JSON.parse(message.value.toString());
await processEvent(event);
}
});
Partitioning Strategy
// Partition CloudEvents by source or subject
// Ensures related events are processed in order
// Partition by orderId for order events
const partitionKey = cloudEvent.data.orderId || cloudEvent.id;
// Kafka producer with partitioning
await producer.send({
topic: 'orders',
messages: [{
key: partitionKey,
value: JSON.stringify(cloudEvent)
}]
});
Event Filtering at Scale
Broker-Based Filtering
# Knative Broker with server-side filtering
apiVersion: eventing.knative.dev/v1
kind: Broker
metadata:
name: events
spec:
# Filter events at broker level
template:
apiVersion: eventing.knative.dev/v1
kind: EventFilter
spec:
type: com.example.order.created
Subscription-Based Filtering
# Knative Subscription with filtering
apiVersion: eventing.knative.dev/v1
kind: Subscription
metadata:
name: order-filtered
spec:
channel:
apiVersion: messaging.knative.dev/v1
kind: Channel
name: events-channel
filter:
attributes:
type: com.example.order.created
source: https://orders.example.com
subscriber:
ref:
apiVersion: serving.knative.dev/v1
kind: Service
name: order-handler
Retry and Dead Letter Handling
Retry Configuration
# Knative Trigger with retry policy
apiVersion: eventing.knative.dev/v1
kind: Trigger
metadata:
name: order-processor
spec:
broker: default
filter:
attributes:
type: com.example.order.created
subscriber:
ref:
apiVersion: serving.knative.dev/v1
kind: Service
name: order-service
retry:
retryOnce: true
backoffPolicy: linear
backoffDelay: "PT1S"
maxRetries: 3
Dead Letter Queue
# Configure dead letter queue for failed events
apiVersion: eventing.knative.dev/v1
kind: Broker
metadata:
name: default
spec:
template:
apiVersion: eventing.knative.dev/v1
kind: BrokerTemplate
spec:
delivery:
deadLetterSink:
ref:
apiVersion: serving.knative.dev/v1
kind: Service
name: dlq-handler
Event Tracking and Auditing
Event Tracking
// Track CloudEvent lifecycle
class EventTracer {
constructor() {
this.traces = new Map();
}
track(event) {
const trace = {
id: event.id,
type: event.type,
source: event.source,
events: [{
timestamp: Date.now(),
action: 'created'
}]
};
this.traces.set(event.id, trace);
}
recordAction(eventId, action, details = {}) {
const trace = this.traces.get(eventId);
if (trace) {
trace.events.push({
timestamp: Date.now(),
action,
...details
});
}
}
getTrace(eventId) {
return this.traces.get(eventId);
}
}
Audit Logging
// Audit logging for CloudEvents
class AuditLogger {
constructor() {
this.queue = [];
this.processing = false;
}
async log(event, context) {
const auditEntry = {
eventId: event.id,
eventType: event.type,
eventSource: event.source,
timestamp: new Date().toISOString(),
userId: context?.userId || 'system',
action: context?.action || 'processed',
dataHash: this.hash(event.data)
};
this.queue.push(auditEntry);
this.processQueue();
}
async processQueue() {
if (this.processing) return;
this.processing = true;
while (this.queue.length > 0) {
try {
const entry = this.queue.shift();
await this.writeEntry(entry);
} catch (error) {
console.error('Audit log write failed:', error);
}
}
this.processing = false;
}
hash(data) {
return crypto.createHash('sha256').update(JSON.stringify(data)).digest('hex');
}
writeEntry(entry) {
// Write to audit log (database, file, etc.)
return auditDatabase.insert(entry);
}
}
Tutorial
Prerequisites
Before working with CloudEvents, ensure you have:
- Kubernetes cluster (v1.20+)
- kubectl (v1.20+)
- Helm (v3.0+)
- Basic understanding of event-driven architectures
- Node.js (v16+ or Python 3.8+) for SDK examples
Check your environment:
kubectl version --client
helm version
node --version
Installation
Using Helm (Recommended)
# Add the Knative repository
helm repo add knative https://charts.knative.dev
helm repo update
# Install Knative Eventing (includes CloudEvents support)
kubectl apply -f https://github.com/knative/eventing/releases/download/knative-v1.12.0/eventing-crds.yaml
kubectl wait --for=condition=Established --timeout=30s -f https://github.com/knative/eventing/releases/download/knative-v1.12.0/eventing-crds.yaml
kubectl apply -f https://github.com/knative/eventing/releases/download/knative-v1.12.0/eventing-core.yaml
kubectl wait --for=condition=Ready --timeout=300s -n knative-eventing --all pods
# Install a broker implementation (e.g., InMemoryChannel)
kubectl apply -f https://github.com/knative/eventing/releases/download/knative-v1.12.0/in-memory-channel.yaml
# Verify installation
kubectl get pods -n knative-eventing
kubectl get brokers.eventing.knative.dev
Using kubectl Direct Apply
# Apply CloudEvents CRDs and components
kubectl apply -f https://github.com/cloudevents/spec/raw/main/cloudevents.yaml
# Install NATS JetStream (alternative event bus)
helm repo add nats https://nats-io.github.io/k8s/helm/
helm repo update
helm install nats nats/nats -n nats-system --create-namespace
# Install Eventing Bus
kubectl apply -f https://github.com/knative/eventing/releases/download/knative-v1.12.0/eventing.yaml
Using Docker for Local Testing
# Run a CloudEvents receiver container
docker run -p 8080:8080 \
-e PORT=8080 \
-e TARGET_EVENT_TYPE=com.example.test \
cloudEvents receiver
# Test locally
curl -X POST http://localhost:8080 \
-H "Content-Type: application/cloudevents+json" \
-H "Ce-Specversion: 1.0" \
-H "Ce-Type: com.example.test" \
-H "Ce-Source: test-client" \
-H "Ce-Id: test-123" \
-d '{"message": "Hello CloudEvents"}'
Basic Configuration
CloudEvents Configuration File
# cloudevents-config.yaml
apiVersion: v1
kind: ConfigMap
metadata:
name: cloudevents-config
namespace: default
data:
config.yaml: |
# CloudEvents configuration
specversion: "1.0"
# Default attributes for events
defaultAttributes:
source: "https://myapp.example.com"
type: "com.example.application.event"
datacontenttype: "application/json"
# Broker configuration
broker:
name: default
namespace: default
# Retry policy
retry:
maxRetries: 3
backoffPolicy: linear
backoffDelay: "PT1S"
# Dead letter queue
deadLetterSink:
uri: "http://dlq-service.default.svc.cluster.local"
Environment Variables
Node.js:
// .env
CLOUD_EVENTS_SOURCE=https://myapp.example.com
CLOUD_EVENTS_TYPE=com.example.order.created
CLOUD_EVENTS_TARGET=http://event-broker.default.svc.cluster.local
CLOUD_EVENTS_RETRY_COUNT=3
CLOUD_EVENTS_TIMEOUT=5000
Python:
# ✅ GOOD — Environment variables for CloudEvents in bash
# .env file equivalent for bash
export CLOUD_EVENTS_SOURCE="https://myapp.example.com"
export CLOUD_EVENTS_TYPE="com.example.order.created"
export CLOUD_EVENTS_TARGET="http://event-broker.default.svc.cluster.local"
export CLOUD_EVENTS_RETRY_COUNT=3
export CLOUD_EVENTS_TIMEOUT=5000
# Source the .env file
set -a
source .env
set +a
# Verify environment
echo "Source: $CLOUD_EVENTS_SOURCE"
echo "Target: $CLOUD_EVENTS_TARGET"
echo "Type: $CLOUD_EVENTS_TYPE"
Kubernetes Secret for Sensitive Configuration
# cloudevents-secrets.yaml
apiVersion: v1
kind: Secret
metadata:
name: cloudevents-secrets
namespace: default
type: Opaque
stringData:
# API keys or tokens for external events
api-key: "your-api-key-here"
webhook-secret: "your-webhook-secret-here"
Usage Examples
Creating CloudEvents in Node.js
// cloudevents-client.js
const { CloudEvent } = require('cloudevents');
// Create a CloudEvent factory
class CloudEventFactory {
constructor(source, defaultTypePrefix = '') {
this.source = source;
this.defaultTypePrefix = defaultTypePrefix;
}
create(type, data, attributes = {}) {
return {
specversion: '1.0',
type: this.defaultTypePrefix + type,
source: this.source,
id: this.generateId(),
time: new Date().toISOString(),
...attributes,
data: data
};
}
generateId() {
return crypto.randomUUID();
}
}
// Usage
const factory = new CloudEventFactory('https://orders.example.com', 'com.example.');
const event = factory.create('order.created', { orderId: '123', total: 99.99 });
console.log('Created CloudEvent:', JSON.stringify(event, null, 2));
Sending CloudEvents via HTTP
// Send CloudEvent to broker
async function sendCloudEvent(event, targetUrl) {
const response = await fetch(targetUrl, {
method: 'POST',
headers: {
'Content-Type': 'application/cloudevents+json',
},
body: JSON.stringify(event)
});
if (!response.ok) {
throw new Error(`Failed to send CloudEvent: ${response.statusText}`);
}
return await response.json();
}
// Usage
const event = factory.create('order.shipped', {
orderId: '123',
trackingNumber: 'TRK-456'
});
await sendCloudEvent(event, 'http://event-broker.default.svc.cluster.local');
Receiving CloudEvents in Node.js
// cloudevents-server.js
const http = require('http');
const server = http.createServer(async (req, res) => {
if (req.method !== 'POST') {
res.writeHead(405);
res.end('Method Not Allowed');
return;
}
// Parse the incoming CloudEvent
const headers = req.headers;
const contentType = headers['content-type'];
let event;
if (contentType === 'application/cloudevents+json') {
// Structured content mode
const body = [];
req.on('data', chunk => body.push(chunk));
req.on('end', () => {
try {
event = JSON.parse(Buffer.concat(body).toString());
handleCloudEvent(event);
res.writeHead(200, { 'Content-Type': 'application/json' });
res.end(JSON.stringify({ status: 'received', id: event.id }));
} catch (err) {
res.writeHead(400);
res.end(JSON.stringify({ error: 'Invalid CloudEvent' }));
}
});
} else if (headers['ce-specversion']) {
// Binary content mode
const body = [];
req.on('data', chunk => body.push(chunk));
req.on('end', () => {
try {
event = {
specversion: headers['ce-specversion'],
type: headers['ce-type'],
source: headers['ce-source'],
id: headers['ce-id'],
time: headers['ce-time'] || new Date().toISOString(),
datacontenttype: headers['content-type'],
data: headers['content-type']?.includes('json')
? JSON.parse(Buffer.concat(body).toString())
: Buffer.concat(body).toString()
};
handleCloudEvent(event);
res.writeHead(200);
res.end();
} catch (err) {
res.writeHead(400);
res.end(JSON.stringify({ error: 'Invalid CloudEvent' }));
}
});
} else {
res.writeHead(400);
res.end(JSON.stringify({ error: 'Not a CloudEvent' }));
}
});
function handleCloudEvent(event) {
console.log(`Received event: ${event.id} of type ${event.type}`);
switch (event.type) {
case 'com.example.order.created':
processOrderCreated(event.data);
break;
case 'com.example.order.shipped':
processOrderShipped(event.data);
break;
default:
console.log('Unknown event type');
}
}
server.listen(8080, () => {
console.log('CloudEvents server running on port 8080');
});
Creating CloudEvents in Python
# ✅ GOOD — Creating CloudEvents with jq and curl
# Generate CloudEvent using jq
EVENT_ID=$(python3 -c "import uuid; print(uuid.uuid4())")
TIMESTAMP=$(date -u +"%Y-%m-%dT%H:%M:%SZ")
EVENT=$(jq -n --arg id "$EVENT_ID" --arg type "com.example.order.created" --arg source "https://orders.example.com" --arg time "$TIMESTAMP" '{
"specversion": "1.0",
"id": $id,
"type": $type,
"source": $source,
"time": $time,
"datacontenttype": "application/json",
"data": {"orderId": "123", "total": 99.99}
}')
echo "$EVENT" | jq .
# Send to broker
curl -s -X POST "http://event-broker.default.svc.cluster.local" -H "Content-Type: application/cloudevents+json" -d "$EVENT" | jq .
Sending CloudEvents via HTTP in Python
# ✅ GOOD — Sending CloudEvents with curl and retry
curl -s -X POST "http://event-broker.default.svc.cluster.local" -H "Content-Type: application/cloudevents+json" -d '{"specversion":"1.0","id":"'$EVENT_ID'","type":"com.example.order.shipped","source":"https://orders.example.com","time":"'$TIMESTAMP'","datacontenttype":"application/json","data":{"orderId":"123","trackingNumber":"TRK-456"}}' | jq .
# With retry logic
for i in $(seq 1 5); do
RESPONSE=$(curl -sf -X POST "http://event-broker.default.svc.cluster.local" -H "Content-Type: application/cloudevents+json" -d "$EVENT" 2>/dev/null)
if [[ $? -eq 0 ]]; then
echo "✅ Event sent: $RESPONSE"
break
fi
echo "⚠️ Retry $i..."
sleep $((2 ** i))
done
Common Operations
Monitoring CloudEvents
Kubernetes Pod Logs:
# Monitor Knative service logs
kubectl logs -l knative-service=order-processor -f
# Monitor event broker logs
kubectl logs -n knative-eventing -l eventing.knative.dev/component=broker-controller -f
# Filter for specific event types
kubectl logs -l knative-service=order-processor | grep "com.example.order.created"
CloudEvents Metrics:
# View CloudEvents metrics
kubectl port-forward -n knative-eventing svc/mt-adapter 8080:8080
# Query Prometheus for CloudEvents metrics
curl -s "http://prometheus-k8s.monitoring.svc.cluster.local:9090/api/v1/query" \
--data-urlencode 'query=event_count_total{type="com.example.order.created"}'
Debugging CloudEvents
Enable Detailed Logging:
# knative-service.yaml with debug logging
apiVersion: serving.knative.dev/v1
kind: Service
metadata:
name: event-processor
spec:
template:
spec:
containers:
- image: event-processor
env:
- name: LOG_LEVEL
value: debug
- name: CE_TYPE
value: com.example.order.created
Test Event Delivery:
# Send test event
kubectl -n default run test-event --rm -i --restart=Never \
--image=alpine/curl:3.18 \
-- curl -s -X POST \
-H "Content-Type: application/cloudevents+json" \
-H "Ce-Specversion: 1.0" \
-H "Ce-Type: com.example.test" \
-H "Ce-Source: test-client" \
-H "Ce-Id: test-123" \
-d '{"message": "test"}' \
http://event-broker.default.svc.cluster.local
Retrying Failed Events
Knative Trigger with Retry Policy:
apiVersion: eventing.knative.dev/v1
kind: Trigger
metadata:
name: order-processor
spec:
broker: default
filter:
attributes:
type: com.example.order.created
subscriber:
ref:
apiVersion: serving.knative.dev/v1
kind: Service
name: order-service
retry:
retryOnce: true
backoffPolicy: linear
backoffDelay: "PT1S"
maxRetries: 3
Best Practices
1. Consistent Event Naming
// ✅ Correct - consistent naming convention
const event = factory.create('order.created', data);
// Type: com.example.order.created
// ❌ Incorrect - inconsistent naming
const badEvent = {
type: 'OrderCreated', // camelCase, no namespace
data: data
};
2. Unique Event IDs
// ✅ Correct - unique ID
const event = factory.create('order.created', data);
// ID: a1b2c3d4-e5f6-7890-abcd-ef1234567890
// ❌ Incorrect - predictable or reused IDs
const badEvent = {
id: '123', // Not unique
...
};
3. Proper Time Formatting
// ✅ Correct - RFC 3339 format
const event = factory.create('order.created', data);
// Time: 2024-01-15T10:30:00Z
// ❌ Incorrect - local time without timezone
const badEvent = {
time: '2024-01-15 10:30:00', // Invalid format
...
};
4. Data Schema Versioning
// Include schema version in events
const event = factory.create('order.created', {
...data,
schemaVersion: '1.0'
});
// Consumer validates schema version
function handleOrderCreated(data) {
if (data.schemaVersion !== '1.0') {
throw new Error(`Unsupported schema version: ${data.schemaVersion}`);
}
// Process event
}
5. Error Events
// Create error events for troubleshooting
function createErrorEvent(originalEvent, error) {
return {
specversion: '1.0',
type: 'com.example.error',
source: originalEvent.source,
id: originalEvent.id + '-error',
time: new Date().toISOString(),
datacontenttype: 'application/json',
data: {
originalEventId: originalEvent.id,
originalType: originalEvent.type,
errorMessage: error.message,
stackTrace: error.stack,
timestamp: new Date().toISOString()
}
};
}
Troubleshooting
Official Documentation
- CloudEvents Specification - GitHub repository
- CloudEvents Website - Project home page
- CloudEvents API Reference - Complete specification
- CloudEvents Use Cases - Real-world applications
Implementations
- CloudEvents SDKs - Official SDKs
- Knative Eventing - Kubernetes-native eventing
- Dapr Pub/Sub - CloudEvents support
- AWS EventBridge - CloudEvents format support
Community Resources
- CNCF CloudEvents - CNCF project page
- CloudEvents Slack - Community discussion
- CloudEvents Mailing List - Announcements and discussions
Learning Resources
- CloudEvents Getting Started - Tutorial
- CloudEvents Examples - Code samples
- CloudEvents Webinars - Video content
Common Issues
-
Deployment Failures
- Check pod logs for errors
- Verify configuration values
- Ensure network connectivity
-
Performance Issues
- Monitor resource usage
- Adjust resource limits
- Check for bottlenecks
-
Configuration Errors
- Validate YAML syntax
- Check required fields
- Verify environment-specific settings
-
Integration Problems
- Verify API compatibility
- Check dependency versions
- Review integration documentation
Getting Help
- Check official documentation
- Search GitHub issues
- Join community channels
- Review logs and metrics
Examples
Basic Configuration
# CloudEvents configuration example
apiVersion: v1
kind: ConfigMap
metadata:
name: cloudevents-config
namespace: default
data:
config.yaml: |
specversion: "1.0"
defaultAttributes:
source: "https://myapp.example.com"
type: "com.example.application.event"
datacontenttype: "application/json"
broker:
name: default
namespace: default
retry:
maxRetries: 3
backoffPolicy: linear
backoffDelay: "PT1S"
Kubernetes Deployment
# Kubernetes deployment for CloudEvents processor
apiVersion: apps/v1
kind: Deployment
metadata:
name: cloudevents-processor
namespace: default
spec:
replicas: 3
selector:
matchLabels:
app: cloudevents-processor
template:
metadata:
labels:
app: cloudevents-processor
spec:
containers:
- name: cloudevents-processor
image: cloudevents-processor:latest
ports:
- containerPort: 8080
env:
- name: CE_SOURCE
value: "https://myapp.example.com"
- name: CE_TYPE
value: "com.example.application.event"
resources:
limits:
memory: "256Mi"
cpu: "500m"
requests:
memory: "128Mi"
cpu: "250m"
Kubernetes Service
# Kubernetes service for CloudEvents processor
apiVersion: v1
kind: Service
metadata:
name: cloudevents-processor
namespace: default
spec:
selector:
app: cloudevents-processor
ports:
- protocol: TCP
port: 80
targetPort: 8080
type: ClusterIP
When to Use
Use this skill when:
- Integrating a CNCF project into Kubernetes infrastructure — You need to configure, deploy, or troubleshoot a cloud-native tool within a cluster
- Designing cloud-native architecture — You are selecting and integrating CNCF tools to solve specific infrastructure challenges
- Resolving operational issues — A CNCF component is misbehaving, underperforming, or needs configuration changes
Core Workflow
-
Assess Requirements — Understand the use case, scale, integration needs, and existing infrastructure. Checkpoint: Document requirements, constraints, and success criteria.
-
Design Architecture — Plan component interactions, data flow, and deployment strategy using cloud-native best practices. Checkpoint: Verify the architecture addresses all requirements and follows CNCF conventions.
-
Implement & Configure — Create manifests, configurations, and deployment scripts. Include resource limits, health checks, and observability hooks. Checkpoint: Validate all YAML against schema and test in a staging environment.
-
Deploy & Monitor — Apply manifests to the cluster, verify component health, and confirm observability is working. Checkpoint: Confirm all pods/services are running, probes passing, and metrics/alerts configured.
Constraints
MUST DO
- Include at least one complete working YAML manifest example
- Note when content is auto-generated vs. manually verified
- Reference relevant CNCF project documentation
MUST NOT DO
- Deploy manifests without testing in a staging environment first
- Use deprecated API versions (e.g., apps/v1beta1)
- Omit resource limits and requests in Kubernetes manifests