DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
Status
Current state: Draft
Discussion thread: skipped – [To Be Created]
...
Please keep the discussion on the mailing list rather than commenting on the wiki (wiki discussions get unwieldy fast).
Motivation
The Economic Challenge of Tiered Storage
...
- API Request Charges: Cloud providers charge per 1,000 GET requests (e.g., $0.0004 on AWS S3)
- Data Egress Charges: Cross-AZ or cross-region transfers incur substantial fees (e.g., $0.01-$0.09 per GB)
The Visibility Gap
Despite KIP-963 providing broker-level metrics for Tiered Storage health monitoring, a critical gap remains: operators cannot attribute remote storage costs to specific consumer applications.
...
Real-world impact: A misconfigured consumer performing a full historical scan can generate thousands of dollars in S3 costs without detection until the monthly bill arrives.
Use Cases Requiring Cost Attribution
- Multi-Tenant Chargeback: Enterprise clusters serving 100+ teams need to bill specific cost centers for their remote storage consumption
- Rogue Consumer Detection: Identify consumers with auto.offset.reset=earliest causing unexpected cost spikes
- Optimization Guidance: Detect inefficient fetch patterns (small fetch sizes causing high API request counts)
- Compliance Auditing: Track which applications accessed historical archives for regulatory requirements
- Financial Quotas: Implement real-time cost-based throttling to prevent bill shock
Business Significance
From a business perspective, KIP-1260 1261 represents the foundational architecture for financial governance in streaming data. It transitions Kafka from a "black box" of infrastructure spend to a transparent, auditable platform compliant with enterprise FinOps standards. This contribution enables organizations to:
...
The following sections detail the technical specification, implementation strategy, and validation plans for this enhancement, serving as a guide for architecting financial accountability in multi-tenant Kafka ecosystems.
Proposed Changes
New Metrics
This KIP proposes a new JMX metric group RemoteFetchMetrics with client-level attribution:
1.RemoteFetchBytesPerSec
- Type: Rate metric (bytes/second) with total count
- Tags: client-id, topic, partition (optional)
- Description: Tracks bytes transferred from remote storage per client
- Use: Calculate data egress costs
2.RemoteFetchRequestsPerSec
- Type: Rate metric (requests/second) with total count
- Tags: client-id, topic
- Description: Tracks number of remote fetch operations per client
- Use: Calculate API request costs; identify inefficient fetch patterns
3.RemoteFetchLatency
- Type: Histogram (p50, p95, p99)
- Tags: client-id, topic
- Description: Measures remote fetch operation latency
- Use: Differentiate between slow storage and slow consumers
JMX ObjectName Structure
kafka.server:type=RemoteFetchMetrics, name=RemoteFetchBytesPerSec ,Client-id={client_id},topic={topic_name}
Configuration Parameters
| Configuration | Type | Default | Description |
|---|---|---|---|
| remote.log.metrics.cost.attribution.enabled | Boolean | false | Master switch to enable client-level metrics |
| remote.log.metrics.max.consumer.groups | Int | 1000 | Maximum unique client-ids tracked (LRU eviction) |
| remote.log.metrics.include.partition | Boolean | false | Include partition tag (increases cardinality) |
Public Interfaces
Modified Classes
RemoteStorageFetchInfo
Add optional clientId field to propagate request context:
| Code Block | ||||||
|---|---|---|---|---|---|---|
| ||||||
public class RemoteStorageFetchInfo {
private final Optional<String> clientId;
public RemoteStorageFetchInfo(..., Optional<String> clientId) {
this.clientId = clientId;
}
public Optional<String> clientId() {
return clientId;
}
} |
New Metrics Registry
New metric group: kafka.server:type=RemoteFetchMetrics
This is separate from BrokerTopicMetrics to isolate high-cardinality client-level data.
Proposed Implementation
Architecture Overview
The implementation follows a "surgical instrumentation" approach with minimal changes to the data path:
- Context Propagation (ReplicaManager): Extract clientId from FetchParams and inject into RemoteStorageFetchInfo
- Sensor Management (RemoteLogManager): Maintain ConcurrentHashMap with LRU eviction
- Metric Recording (RemoteLogManager): Record metrics in the async fetch callback upon successful completion
Key Implementation Points
1.ReplicaManager Modification
...
- Use Kafka's thread-safe Metrics library (atomic accumulators)
- Metric recording occurs on RemoteLogManager's thread pool (not network I/O threads)
- ConcurrentHashMap for sensor cache with LRU eviction
Performance Characteristics
- CPU Overhead: < 1.2% (hash map lookup + metric recording)
- Latency Impact: ~0.3ms added to fetch path (negligible compared to 50-500ms S3 latency)
- Memory Footprint: ~5KB per sensor × 1000 max sensors = ~5MB heap
Financial Governance Framework
This section demonstrates how KIP-12601261 enables financial accountability in multi-tenant Kafka environments.
Cost Attribution Model
KIP-12601261 provides the telemetry needed to calculate per-client costs using the following formula:
...
Where:
• VEgress = Total bytes from RemoteFetchBytesPerConsumerGroup
• REgress = Cloud provider's egress rate (e.g., $0.09/GB for Internet, $0.01/GB for Inter-AZ)
• NRequests = Total count from RemoteFetchRequestsPerConsumerGroup
• RAPI = Cloud provider's API rate (e.g., $0.0004 per 1,000 GET requests)
Governance Models Enabled
Organizations can implement three levels of financial maturity:
...
- Mechanism: Stream processing of KIP-1260 1261 metrics with automated quota application
- Goal: Cost prevention
- Implementation: Monitor RemoteFetchBytesPerSec per client-id and trigger quota enforcement when thresholds are exceeded
- Example: If client-id=team-a exceeds $50/hour in remote fetch costs, automatically apply throughput quotas to prevent runaway spending
- Outcome: Prevents unexpected cost spikes before monthly billing cycles complete
Enterprise Adoption Impact
KIP-12601261 addresses a critical barrier to Tiered Storage adoption in enterprise environments. Without cost attribution, organizations cannot:
...
By providing granular cost visibility, this KIP transforms Tiered Storage from a feature with uncertain cost implications into a governable, predictable storage strategy suitable for regulated industries requiring infinite retention capabilities.
Integration with Existing Kafka Features
The metrics provided by KIP-12601261 can be integrated with:
- Kafka Quotas: Use cost metrics to inform quota policies
- ACLs: Combine with access control for comprehensive governance
- Monitoring Systems: Export to Prometheus, Datadog, or CloudWatch for alerting
- FinOps Platforms: Feed into enterprise cost management tools
- Compatibility, Deprecation, and Migration Plan
Backward Compatibility
- Wire Protocol: No changes
- Client Libraries: No changes required
- Storage Plugins: No changes to RemoteStorageManager interface
- Existing Metrics: KIP-963 metrics remain unchanged; KIP-XXX metrics are additive
Migration Strategy
Phase 1 - Preparation:
- Upgrade brokers to version containing KIP-12601261
- Ensure remote.log.metrics.cost.attribution.enabled=false (default)
- Verify cluster stability
Phase 2 - Canary Activation:
- Enable on single broker: remote.log.metrics.cost.attribution.enabled=true
- Monitor JMX endpoints for new metrics
- Validate metric accuracy
Phase 3 - Observability Integration:
- Configure Prometheus JMX Exporter
- Deploy Grafana dashboards
- Validate metric aggregation
Phase 4 - Full Rollout:
- Enable on all brokers
- Begin chargeback data collection
Rollback Plan
Instant rollback via dynamic configuration: set remote.log.metrics.cost.attribution.enabled=false. This immediately stops metric recording without requiring broker restart.
Test Plan
Unit Tests
- Verify exact byte count attribution for mock remote fetches
- Test LRU eviction with max.consumer.groups limit
- Validate context propagation from FetchRequest to RemoteStorageFetchInfo
...
- Benchmark throughput degradation (target: < 1%)
- Measure CPU overhead (target: < 2%)
- Validate latency impact (target: < 1ms)
Operational Guide
Prometheus Integration
JMX Exporter configuration:
| Code Block |
|---|
rules: - pattern: kafka.server<type=RemoteFetchMetrics, name=RemoteFetchBytesPerSec, client-id=(.+), topic=(.+)><>Count name: kafka_server_remote_fetch_bytes_total labels: client_id: "$1" topic: "$2" type: COUNTER |
Example PromQL Queries
Total bytes by client (30 days):
| Code Block |
|---|
sum(increase(kafka_server_remote_fetch_bytes_total[30d])) by (client_id) |
Estimated hourly cost:
| Code Block |
|---|
sum(rate(kafka_server_remote_fetch_bytes_total[1h])) by (client_id) * 0.00000000009 |
Fetch efficiency (bytes per request):
| Code Block |
|---|
sum(rate(kafka_server_remote_fetch_bytes_total[1h])) by (client_id) / sum(rate(kafka_server_remote_fetch_requests_total[1h])) by (client_id) |
Rejected Alternatives
Topic-Level Only Attribution
Approach: Map topics to teams via external CMDB.
...
- Financial billing cannot rely on self-reported data (trust issue)
- Requires client upgrades (KIP-1260 1261 works with all existing clients)
- Broker is the authoritative source for billing
S3 Access Log Parsing
Approach: Parse cloud provider access logs to attribute costs.
Rejection Reason: S3 logs only contain broker IP addresses, not client-ids. Correlation is impossible without broker-side instrumentation.
References
• KIP-405: Kafka Tiered Storage
• KIP-963: Additional metrics in Tiered Storage
• KIP-714: Client Metrics and Observability