DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
...
| Code Block | ||
|---|---|---|
| ||
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 getVersion() {
return VERSION;
}
public static String getCommitId() {
return COMMIT_ID;
}
public static String getClientId() {
return CLIENT_ID;
}
public static synchronized void registerAppInfo(String prefix, String id, Metrics metrics, long nowMs) {
try {
ObjectName name = new ObjectName(prefix + ":type=app-info,id=" + Sanitizer.jmxSanitize(id));
MBeanServer server = ManagementFactory.getPlatformMBeanServer();
if (server.isRegistered(name)) {
log.info("The mbean of App info: [{}], id: [{}] already exists, so skipping a new mbean creation.", prefix, id);
return;
}
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);
}
}
public static synchronized void unregisterAppInfo(String prefix, String id, Metrics metrics) {
MBeanServer server = ManagementFactory.getPlatformMBeanServer();
try {
ObjectName name = new ObjectName(prefix + ":type=app-info,id=" + Sanitizer.jmxSanitize(id));
if (server.isRegistered(name))
server.unregisterMBean(name);
unregisterMetrics(metrics);
} catch (JMException e) {
log.warn("Error unregistering AppInfo mbean", e);
} finally {
log.info("App info {} for {} unregistered", prefix, id);
}
}
private static MetricName metricName(Metrics metrics, String name) {
return metrics.metricName(name, "app-info", "Metric indicating " + name);
}
private static void registerMetrics(Metrics metrics, AppInfoMBean appInfo) {
if (metrics != null) {
metrics.addMetric(metricName(metrics, "version"), new ImmutableValue<>(appInfo.getVersion()));
metrics.addMetric(metricName(metrics, "commit-id"), new ImmutableValue<>(appInfo.getCommitId()));
metrics.addMetric(metricName(metrics, "start-time-ms"), new ImmutableValue<>(appInfo.getStartTimeMs()));
if (appInfo instanceof AppInfo) {
metrics.addMetric(metricName(metrics, "client-id"), new ImmutableValue<>(appInfo.getClientId()));
}
}
}
private static void unregisterMetrics(Metrics metrics) {
if (metrics != null) {
metrics.removeMetric(metricName(metrics, "version"));
metrics.removeMetric(metricName(metrics, "commit-id"));
metrics.removeMetric(metricName(metrics, "start-time-ms"));
metrics.removeMetric(metricName(metrics, "client-id"));
}
}
public interface AppInfoMBean {
String getVersion();
String getCommitId();
String getClientId();
Long getStartTimeMs();
}
@Deprecated(since = "4.2")
public static class DeprecatedAppInfo implements AppInfoMBean {
private final Long startTimeMs;
public DeprecatedAppInfo(long startTimeMs) {
this.startTimeMs = startTimeMs;
log.info("Kafka version: {}", AppInfoParser.getVersion());
log.info("Kafka commitId: {}", AppInfoParser.getCommitId());
log.info("Kafka startTimeMs: {}", startTimeMs);
}
@Override
public String getVersion() {
return AppInfoParser.getVersion();
}
@Override
public String getCommitId() {
return AppInfoParser.getCommitId();
}
@Override
public String getClientId() {
throw new UnsupportedOperationException("client-id is not implemented in DeprecatedAppInfo");
}
@Override
public Long getStartTimeMs() {
return startTimeMs;
}
}
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());
}
@Override
public String getVersion() {
return AppInfoParser.getVersion();
}
@Override
public String getCommitId() {
return AppInfoParser.getCommitId();
}
@Override
public String getClientId() {
return AppInfoParser.getClientId();
}
@Override
public Long getStartTimeMs() {
return startTimeMs;
}
}
static class ImmutableValue<T> implements Gauge<T> {
private final T value;
public ImmutableValue(T value) {
this.value = value;
}
@Override
public T value(MetricConfig config, long now) {
return value;
}
}
} |
Compatibility, Deprecation, and Migration Plan
...