Versions Compared

Key

  • This line was added.
  • This line was removed.
  • Formatting was changed.


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

  1. Multi-Tenant Chargeback: Enterprise clusters serving 100+ teams need to bill specific cost centers for their remote storage consumption
  2. Rogue Consumer Detection: Identify consumers with auto.offset.reset=earliest causing unexpected cost spikes
  3. Optimization Guidance: Detect inefficient fetch patterns (small fetch sizes causing high API request counts)
  4. Compliance Auditing: Track which applications accessed historical archives for regulatory requirements
  5. 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

ConfigurationTypeDefaultDescription
remote.log.metrics.cost.attribution.enabledBooleanfalseMaster switch to enable client-level metrics
remote.log.metrics.max.consumer.groupsInt1000Maximum unique client-ids tracked (LRU eviction)
remote.log.metrics.include.partitionBooleanfalseInclude partition tag (increases cardinality)


Public Interfaces

Modified Classes

RemoteStorageFetchInfo


Add optional clientId field to propagate request context:

Code Block
languagejava
themeEclipse
titlecode
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:

  1. Upgrade brokers to version containing KIP-12601261
  2. Ensure remote.log.metrics.cost.attribution.enabled=false (default)
  3. Verify cluster stability

Phase 2 - Canary Activation:

  1. Enable on single broker: remote.log.metrics.cost.attribution.enabled=true
  2. Monitor JMX endpoints for new metrics
  3. Validate metric accuracy

Phase 3 - Observability Integration:

  1. Configure Prometheus JMX Exporter
  2. Deploy Grafana dashboards
  3. Validate metric aggregation

Phase 4 - Full Rollout:

  1. Enable on all brokers
  2. 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