Versions Compared

Key

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

Table of Contents

This page is meant as a template for writing a KIP. To create a KIP choose Tools->Copy on this page and modify with your content and replace the heading with the next KIP number and a description of your issue. Replace anything in italics with your own description.

Status

Current state:  DraftUnder Discussion

Discussion thread: here [Change the link from the KIP proposal email archive to your own email thread]

JIRA: KAFKA-20297

Please keep the discussion on the mailing list rather than commenting on the wiki (wiki discussions get unwieldy fast).

Motivation

Current state

The The Time interface and Timer class in org.apache.kafka.common.utils is are currently exposed through public APIs but not officially designated as public. It serves multiple responsibilities:

  1. Wall clock time (milliseconds())
  2. Monotonic high-resolution timing (nanoseconds(), hiResClockMs())
  3. Thread coordination (sleep(), waitObject(), waitForFuture())
  4. Timer factory (timer() methods)

Classes like KafkaStreams accept Time parameters, making it part of the user-facing API, yet it is not included in the public Javadocs or officially supported. So, before making it official, we should address design issues identified by the community.

Problems with the current Time interface:

  • Conflates unrelated responsibilities - mixes time reading, thread coordination, and timer creation
  • Testing pain - mocking Time requires implementing 8 methods even if the code only uses 1
  • Unclear dependencies - The Time parameter doesn't indicate what functionality is needed
  • Usage imbalance - Analysis shows milliseconds() has 311 uses, while waitForFuture() has 0 uses

Community feedback during KIP-1247 discussion identified Time as combining unrelated responsibilities and being primarily used for testing. This KIP addresses those concerns before making it officially public.

Public Interfaces

Briefly list any new interfaces that will be introduced as part of this proposal or any existing interfaces that will be removed or changed. The purpose of this section is to concisely call out the public contract that will come along with this feature.

A public interface is any change to the following:

  • Binary log format

  • The network protocol and api behavior

  • Any class in the public packages under clientsConfiguration, especially client configuration

    • org/apache/kafka/common/serialization

    • org/apache/kafka/common

    • org/apache/kafka/common/errors

    • org/apache/kafka/clients/producer

    • org/apache/kafka/clients/consumer (eventually, once stable)

  • Monitoring

  • Command line tools and arguments

  • Anything else that will likely break existing users in some way when they upgrade

Proposed Changes

Describe the new thing you want to do in appropriate detail. This may be fairly extensive and have large subsections of its own. Or it may be a few sentences. Use judgement based on the scope of the change.

Compatibility, Deprecation, and Migration Plan

  • What impact (if any) will there be on existing users?
  • If we are changing behavior how will we phase out the older behavior?
  • If we need special migration tools, describe them here.
  • When will we remove the existing behavior?

Test Plan

Describe in few sentences how the KIP will be tested. We are mostly interested in system tests (since unit-tests are specific to implementation details). How will we know that the implementation works as expected? How will we know nothing broke?

Rejected Alternatives

Following KIP-1247 (Make Bytes public), these should be made officially public.

Code Block
languagejava
// multiple KafkaStreams constructors (test-only usage)
public KafkaStreams(Topology topology, Properties props, Time time)

// couple of examples from multiple "org.apache.kafka.common.metrics.Metrics.java" constructors
public Metrics(Time time) {}
public Metrics(MetricConfig defaultConfig, Time time) {}

// in package "org.apache.kafka.common.security.oauthbearer"
public JwtBearerJwtRetriever(Time time) {}
public ClientCredentialsJwtRetriever(Time time) {}


However, KafkaStreams exposes constructors accepting Time parameters that are only useful for testing, not production use. As noted in the KIP-1247 discussion, these test-only constructors should be deprecated before making Time the official public API.

Timer is returned by Time interface methods and users might interact with Timer's public methods. If Time becomes officially public, Timer must also be public as it is part of Time's contract.

Public Interfaces

New public APIs (4.4):

  • org.apache.kafka.common.utils.Time (officially public)
  • org.apache.kafka.common.utils.Timer (officially public)
  • org.apache.kafka.streams.test.KafkaStreamsMock (test utility)

Deprecated (4.4, removed in 5.0):

  • KafkaStreams(Topology, Properties, Time) constructor
  • KafkaStreams(Topology, Properties, KafkaClientSupplier, Time) constructor

These are only used for testing and pollute the public API.

Proposed Changes

Make Time and Timer officially public (4.4):

Include Time and Timer in public API Javadocs with enhanced documentation.

Deprecate KafkaStreams test constructors (4.4):

Deprecate the following constructors (removal in 5.0):

  • KafkaStreams(Topology, Properties, Time)
  • KafkaStreams(Topology, Properties, KafkaClientSupplier, Time)

Provide Test Utility (4.4):

Create org.apache.kafka.streams.test.KafkaStreamsMock with static factory methods for creating KafkaStreams instances with custom Time implementations.

Example:

Code Block
languagejava
MockTime mockTime = new MockTime();
KafkaStreams streams = KafkaStreamsMock.create(topology, props, mockTime);

Compatibility, Deprecation, and Migration Plan

  • Version 4.4: Time and Timer become officially public. KafkaStreams Time constructors will be deprecated but functional.
  • Versions 4.4-4.9: Deprecation period for test migration.
  • Version 5.0: Deprecated constructors removed (become package-private, accessible via KafkaStreamsMock).

Test migration:

Code Block
languagejava
// Before
KafkaStreams streams = new KafkaStreams(topology, props, mockTime);

// After
KafkaStreams streams = KafkaStreamsMock.create(topology, props, mockTime);

No changes required for production code.

Test Plan

This KIP primarily involves API designation and deprecation rather than functional changes. Testing will focus on:

  • Javadoc generation: Verify Time and Timer appear correctly in public API documentation
  • KafkaStreamsMock Utility: Unit tests confirming the utility creates functional KafkaStreams instances with custom Time implementations
  • Backward Compatibility: Existing tests using deprecated constructors continue to pass unchanged
  • Deprecation Warnings: Build verification that deprecated constructors produce appropriate warnings

No new system tests required. Time and Timer functionality remains unchanged, only their API designation and test access patterns are updated.

Rejected Alternatives

Keep test constructors public: Continues to pollute the public API with test-only functionalityIf there are alternative ways of accomplishing the same thing, what were they? The purpose of this section is to motivate why the design is the way it is and not some other way.