DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
Status
Current state: Under Discussion
Discussion thread: here [Change the link from the KIP proposal email archive to your own email thread]
JIRA: KAFKA-15186 - Getting issue details... STATUS
Motivation
All Kafka component register AppInfo metrics to track the application start time or commit-id etc. These metrics are useful for monitoring and debugging. However, the AppInfo doesn't provide client-id, which is an important information for custom metrics reporter.
As the AppInfoParser class registers a JMX MBean with the provided client-id but when it adds metrics to the Metrics registry the client-id is not included.
Public Interfaces
org.apache.kafka.common.utils.AppInfoParser.AppInfoMBean
Proposed Changes
1) Update the AppInfoMBean interface add new method getClientId()
public interface AppInfoMBean {
String getVersion();
String getCommitId();
Long getStartTimeMs();
String getClientId();
}
2) Deprecated old AppInfoMBean implementation and mark it deprecated
@Deprecated(since = "4.2")
public static class DeprecatedAppInfo implements AppInfoMBean {
}
3) Add new AppInfoMBeanimplementation
public static class AppInfo implements AppInfoMBean {
private final Long startTimeMs;
public AppInfo(long startTimeMs) {
this.startTimeMs = startTimeMs;
log.info("Kafka version: {}", AppInfoParser.getVersion());
log.info("Kafka commitId: {}", AppInfoParser.getCommitId());
log.info("Kafka startTimeMs: {}", startTimeMs);
log.info("Kafka client id: {}", AppInfoParser.getClientId());
}
// skip...
@Override
public String getClientId() {
return AppInfoParser.getClientId();
}
}
4) When AppInfoParser register MBean, we will register the new one and the deprecated one
public class AppInfoParser {
private static final Logger log = LoggerFactory.getLogger(AppInfoParser.class);
private static final String VERSION;
private static final String COMMIT_ID;
private static final String CLIENT_ID;
protected static final String DEFAULT_VALUE = "unknown";
static {
Properties props = new Properties();
try (InputStream resourceStream = AppInfoParser.class.getResourceAsStream("/kafka/kafka-version.properties")) {
props.load(resourceStream);
} catch (Exception e) {
log.warn("Error while loading kafka-version.properties: {}", e.getMessage());
}
VERSION = props.getProperty("version", DEFAULT_VALUE).trim();
COMMIT_ID = props.getProperty("commitId", DEFAULT_VALUE).trim();
CLIENT_ID = props.getProperty("clientId", DEFAULT_VALUE).trim();
}
public static String getClientId() {
return CLIENT_ID;
}
public static synchronized void registerAppInfo(String prefix, String id, Metrics metrics, long nowMs) {
try {
// skip...
DeprecatedAppInfo deprecatedMBean = new DeprecatedAppInfo(nowMs);
AppInfo mBean = new AppInfo(nowMs);
server.registerMBean(deprecatedMBean, name);
server.registerMBean(mBean, name);
registerMetrics(metrics, mBean); // prefix will be added later by JmxReporter
} catch (JMException e) {
log.warn("Error registering AppInfo mbean", e);
}
}
}
Compatibility, Deprecation, and Migration Plan
Since this is a newly introduced behavior, there are no compatibility concerns.
Test Plan
change unit test or Integration test to verify the new method.
Rejected Alternatives
n/a