Skip to content

Commit cac4d1f

Browse files
authored
MINOR: Remove clients test-fixtures shadow JAR rewire hack (#22699)
See the discussion: #22201 (comment) Reviewers: Chia-Ping Tsai <chia7712@gmail.com>
1 parent 624ca39 commit cac4d1f

10 files changed

Lines changed: 400 additions & 517 deletions

File tree

build.gradle

Lines changed: 0 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -1982,28 +1982,6 @@ project(':clients') {
19821982
shadowed
19831983
}
19841984

1985-
// Rewire test fixtures dependencies to avoid depending on the shadow JAR:
1986-
// java-test-fixtures creates dependencies testImplementation -> testFixturesApi -> (main artifact)
1987-
// It prefers the main artifact over the raw main classpath so that dependents on fixtures get the same
1988-
// packaged artifact as production consumers. In this module the main artifact is the shadow JAR,
1989-
// so this creates a test dependency on createVersionFile breaking AppInfoParserTest, and worse:
1990-
// test code is compiled against the original packages, but the shadow JAR contains relocated bytecode,
1991-
// causing runtime errors for various tests. To fix this, we rewire the test fixtures dependency to the raw
1992-
// source set output, so tests continue to use compiled classes and original dependencies instead of the shadow JAR.
1993-
// https://github.com/gradle/gradle/blob/v9.4.1/platforms/jvm/plugins-jvm-test-fixtures/src/main/java/org/gradle/api/plugins/JavaTestFixturesPlugin.java#L89-L90
1994-
afterEvaluate {
1995-
configurations.testFixturesApi.dependencies.removeIf { dep ->
1996-
dep instanceof ProjectDependency && dep.name == project.name
1997-
}
1998-
dependencies {
1999-
testFixturesApi files(sourceSets.main.output.classesDirs)
2000-
testFixturesApi files(sourceSets.main.output.resourcesDir)
2001-
}
2002-
tasks.named('compileTestFixturesJava') {
2003-
dependsOn tasks.named('classes')
2004-
}
2005-
}
2006-
20071985
dependencies {
20081986
implementation libs.zstd
20091987
implementation libs.lz4

clients/src/main/java/org/apache/kafka/common/telemetry/internals/ClientTelemetryProvider.java

Lines changed: 15 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -23,9 +23,9 @@
2323
import org.apache.kafka.common.metrics.MetricsContext;
2424

2525
import java.util.HashMap;
26+
import java.util.LinkedHashMap;
2627
import java.util.List;
2728
import java.util.Map;
28-
import java.util.stream.Collectors;
2929

3030
import io.opentelemetry.proto.common.v1.AnyValue;
3131
import io.opentelemetry.proto.common.v1.KeyValue;
@@ -117,8 +117,7 @@ synchronized void contextChange(MetricsContext metricsContext) {
117117
*/
118118
synchronized void updateLabels(Map<String, String> labels) {
119119
final Resource.Builder resourceBuilder = resource.toBuilder();
120-
Map<String, String> finalLabels = resource.getAttributesList().stream().collect(Collectors.toMap(
121-
KeyValue::getKey, kv -> kv.getValue().getStringValue()));
120+
Map<String, String> finalLabels = resourceLabels();
122121
finalLabels.putAll(labels);
123122

124123
resourceBuilder.clearAttributes();
@@ -135,6 +134,19 @@ Resource resource() {
135134
return resource;
136135
}
137136

137+
/**
138+
* The resource attributes as plain labels.
139+
*
140+
* @return resource attributes keyed by label name, preserving insertion order.
141+
*/
142+
synchronized Map<String, String> resourceLabels() {
143+
Map<String, String> labels = new LinkedHashMap<>();
144+
for (KeyValue attribute : resource.getAttributesList()) {
145+
labels.put(attribute.getKey(), attribute.getValue().getStringValue());
146+
}
147+
return labels;
148+
}
149+
138150
/**
139151
* Domain of the active provider i.e. specifies prefix to the metrics.
140152
*

clients/src/main/java/org/apache/kafka/common/telemetry/internals/ClientTelemetryUtils.java

Lines changed: 26 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,9 @@
4343
import java.util.function.Predicate;
4444

4545
import io.opentelemetry.proto.metrics.v1.MetricsData;
46+
import io.opentelemetry.proto.metrics.v1.ResourceMetrics;
47+
import io.opentelemetry.proto.metrics.v1.ScopeMetrics;
48+
import io.opentelemetry.proto.resource.v1.Resource;
4649

4750
public class ClientTelemetryUtils {
4851

@@ -214,6 +217,29 @@ public static ByteBuffer compress(MetricsData metrics, CompressionType compressi
214217
}
215218
}
216219

220+
/**
221+
* Assembles and compresses a {@code MetricsData} payload for the given metrics. The io.opentelemetry
222+
* types are relocated in the shaded clients jar, so keeping this assembly in main code lets test code
223+
* build and compress a payload without referencing io.opentelemetry directly. The layout mirrors the
224+
* payload produced by {@link ClientTelemetryReporter} (one resource metric per single point metric).
225+
*
226+
* @param metrics the metrics to assemble into a {@code MetricsData} payload
227+
* @param compressionType the compression to apply
228+
* @return the compressed {@code MetricsData} payload
229+
*/
230+
public static ByteBuffer compressMetrics(List<SinglePointMetric> metrics, CompressionType compressionType) throws IOException {
231+
MetricsData.Builder builder = MetricsData.newBuilder();
232+
for (SinglePointMetric metric : metrics) {
233+
builder.addResourceMetrics(ResourceMetrics.newBuilder()
234+
.setResource(Resource.newBuilder().build())
235+
.addScopeMetrics(ScopeMetrics.newBuilder()
236+
.addMetrics(metric.builder())
237+
.build())
238+
.build());
239+
}
240+
return compress(builder.build(), compressionType);
241+
}
242+
217243
public static ByteBuffer decompress(ByteBuffer metrics, CompressionType compressionType, int maxDecompressedBytes) {
218244
Compression compression = Compression.of(compressionType).build();
219245
try (InputStream in = compression.wrapForInput(metrics, RecordBatch.CURRENT_MAGIC_VALUE, BufferSupplier.create());
@@ -235,14 +261,6 @@ public static ByteBuffer decompress(ByteBuffer metrics, CompressionType compress
235261
}
236262
}
237263

238-
public static MetricsData deserializeMetricsData(ByteBuffer serializedMetricsData) {
239-
try {
240-
return MetricsData.parseFrom(serializedMetricsData);
241-
} catch (IOException e) {
242-
throw new KafkaException("Unable to parse MetricsData payload", e);
243-
}
244-
}
245-
246264
public static Uuid fetchClientInstanceId(ClientTelemetryReporter clientTelemetryReporter, Duration timeout) {
247265
if (timeout.isNegative()) {
248266
throw new IllegalArgumentException("The timeout cannot be negative.");

clients/src/main/java/org/apache/kafka/common/telemetry/internals/SinglePointMetric.java

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717
package org.apache.kafka.common.telemetry.internals;
1818

1919
import java.time.Instant;
20+
import java.util.LinkedHashMap;
2021
import java.util.Map;
2122
import java.util.Set;
2223
import java.util.concurrent.TimeUnit;
@@ -50,6 +51,71 @@ public Metric.Builder builder() {
5051
return metricBuilder;
5152
}
5253

54+
// Visible for testing
55+
boolean hasSum() {
56+
return metricBuilder.hasSum();
57+
}
58+
59+
// Visible for testing
60+
boolean hasGauge() {
61+
return metricBuilder.hasGauge();
62+
}
63+
64+
// Visible for testing
65+
boolean isMonotonic() {
66+
return metricBuilder.getSum().getIsMonotonic();
67+
}
68+
69+
// Visible for testing
70+
boolean isDeltaTemporality() {
71+
return metricBuilder.getSum().getAggregationTemporality() == AggregationTemporality.AGGREGATION_TEMPORALITY_DELTA;
72+
}
73+
74+
// Visible for testing
75+
int dataPointsCount() {
76+
return metricBuilder.hasSum()
77+
? metricBuilder.getSum().getDataPointsCount()
78+
: metricBuilder.getGauge().getDataPointsCount();
79+
}
80+
81+
// Visible for testing
82+
long timeUnixNano() {
83+
return dataPoint().getTimeUnixNano();
84+
}
85+
86+
// Visible for testing
87+
long startTimeUnixNano() {
88+
return dataPoint().getStartTimeUnixNano();
89+
}
90+
91+
// Visible for testing
92+
double doubleValue() {
93+
return dataPoint().getAsDouble();
94+
}
95+
96+
// Visible for testing
97+
long longValue() {
98+
return dataPoint().getAsInt();
99+
}
100+
101+
// Visible for testing
102+
int attributesCount() {
103+
return dataPoint().getAttributesCount();
104+
}
105+
106+
// Visible for testing
107+
Map<String, String> attributes() {
108+
Map<String, String> attributes = new LinkedHashMap<>();
109+
for (KeyValue attribute : dataPoint().getAttributesList()) {
110+
attributes.put(attribute.getKey(), attribute.getValue().getStringValue());
111+
}
112+
return attributes;
113+
}
114+
115+
private NumberDataPoint dataPoint() {
116+
return metricBuilder.hasSum() ? metricBuilder.getSum().getDataPoints(0) : metricBuilder.getGauge().getDataPoints(0);
117+
}
118+
53119
/*
54120
Methods to construct gauge metric type.
55121
*/

clients/src/test/java/org/apache/kafka/common/requests/PushTelemetryRequestTest.java

Lines changed: 21 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -36,12 +36,6 @@
3636
import java.util.Collections;
3737
import java.util.List;
3838

39-
import io.opentelemetry.proto.metrics.v1.Metric;
40-
import io.opentelemetry.proto.metrics.v1.MetricsData;
41-
import io.opentelemetry.proto.metrics.v1.ResourceMetrics;
42-
import io.opentelemetry.proto.metrics.v1.ScopeMetrics;
43-
import io.opentelemetry.proto.resource.v1.Resource;
44-
4539
import static org.junit.jupiter.api.Assertions.assertArrayEquals;
4640
import static org.junit.jupiter.api.Assertions.assertEquals;
4741
import static org.junit.jupiter.api.Assertions.assertNotNull;
@@ -59,24 +53,22 @@ public void testGetErrorResponse() {
5953
@ParameterizedTest
6054
@EnumSource(CompressionType.class)
6155
public void testMetricsDataCompression(CompressionType compressionType) throws IOException {
62-
MetricsData metricsData = getMetricsData();
63-
PushTelemetryRequest req = getPushTelemetryRequest(metricsData, compressionType);
56+
List<SinglePointMetric> metrics = sampleMetrics();
57+
byte[] raw = Utils.toArray(ClientTelemetryUtils.compressMetrics(metrics, CompressionType.NONE));
58+
PushTelemetryRequest req = getPushTelemetryRequest(metrics, raw, compressionType);
6459

6560
ByteBuffer receivedMetricsBuffer = req.metricsData(1024 * 1024);
6661
assertNotNull(receivedMetricsBuffer);
6762
assertTrue(receivedMetricsBuffer.capacity() > 0);
68-
69-
MetricsData receivedData = ClientTelemetryUtils.deserializeMetricsData(receivedMetricsBuffer);
70-
assertEquals(metricsData, receivedData);
63+
assertArrayEquals(raw, Utils.toArray(receivedMetricsBuffer));
7164
}
7265

73-
private PushTelemetryRequest getPushTelemetryRequest(MetricsData metricsData, CompressionType compressionType) throws IOException {
74-
ByteBuffer compressedData = ClientTelemetryUtils.compress(metricsData, compressionType);
75-
byte[] data = metricsData.toByteArray();
66+
private PushTelemetryRequest getPushTelemetryRequest(List<SinglePointMetric> metrics, byte[] raw, CompressionType compressionType) throws IOException {
67+
ByteBuffer compressedData = ClientTelemetryUtils.compressMetrics(metrics, compressionType);
7668
if (compressionType != CompressionType.NONE) {
77-
assertTrue(compressedData.limit() < data.length);
69+
assertTrue(compressedData.limit() < raw.length);
7870
} else {
79-
assertArrayEquals(Utils.toArray(compressedData), data);
71+
assertArrayEquals(Utils.toArray(compressedData), raw);
8072
}
8173

8274
return new PushTelemetryRequest.Builder(
@@ -85,36 +77,19 @@ private PushTelemetryRequest getPushTelemetryRequest(MetricsData metricsData, Co
8577
.setCompressionType(compressionType.id)).build();
8678
}
8779

88-
private MetricsData getMetricsData() {
89-
List<Metric> metricsList = new ArrayList<>();
90-
metricsList.add(SinglePointMetric.sum(
91-
new MetricKey("metricName"), 1.0, true, Instant.now(), null, Collections.emptySet())
92-
.builder().build());
93-
metricsList.add(SinglePointMetric.sum(
94-
new MetricKey("metricName1"), 100.0, false, Instant.now(), Instant.now(), Collections.emptySet())
95-
.builder().build());
96-
metricsList.add(SinglePointMetric.deltaSum(
97-
new MetricKey("metricName2"), 1.0, true, Instant.now(), Instant.now(), Collections.emptySet())
98-
.builder().build());
99-
metricsList.add(SinglePointMetric.gauge(
100-
new MetricKey("metricName3"), 1.0, Instant.now(), Collections.emptySet())
101-
.builder().build());
102-
metricsList.add(SinglePointMetric.gauge(
103-
new MetricKey("metricName4"), Long.valueOf(100), Instant.now(), Collections.emptySet())
104-
.builder().build());
105-
106-
MetricsData.Builder builder = MetricsData.newBuilder();
107-
for (Metric metric : metricsList) {
108-
ResourceMetrics rm = ResourceMetrics.newBuilder()
109-
.setResource(Resource.newBuilder().build())
110-
.addScopeMetrics(ScopeMetrics.newBuilder()
111-
.addMetrics(metric)
112-
.build()
113-
).build();
114-
builder.addResourceMetrics(rm);
115-
}
116-
117-
return builder.build();
80+
private List<SinglePointMetric> sampleMetrics() {
81+
List<SinglePointMetric> metrics = new ArrayList<>();
82+
metrics.add(SinglePointMetric.sum(
83+
new MetricKey("metricName"), 1.0, true, Instant.now(), null, Collections.emptySet()));
84+
metrics.add(SinglePointMetric.sum(
85+
new MetricKey("metricName1"), 100.0, false, Instant.now(), Instant.now(), Collections.emptySet()));
86+
metrics.add(SinglePointMetric.deltaSum(
87+
new MetricKey("metricName2"), 1.0, true, Instant.now(), Instant.now(), Collections.emptySet()));
88+
metrics.add(SinglePointMetric.gauge(
89+
new MetricKey("metricName3"), 1.0, Instant.now(), Collections.emptySet()));
90+
metrics.add(SinglePointMetric.gauge(
91+
new MetricKey("metricName4"), Long.valueOf(100), Instant.now(), Collections.emptySet()));
92+
return metrics;
11893
}
11994

12095
}

0 commit comments

Comments
 (0)