-
Notifications
You must be signed in to change notification settings - Fork 15.4k
KAFKA-15778 & KAFKA-15779: Implement metrics manager (KIP-714) #14699
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from 6 commits
d2499d0
26db98e
31a1e23
fb4251f
d4f70e1
84038aa
8f11a0e
c5e2098
177edb4
83963a0
42fe843
05b1100
7f69846
8914649
26fd7bd
eb8463d
8afb5e1
f9efcfb
d4ea0a5
b92d5fa
46ad39f
deef326
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,28 @@ | ||
| /* | ||
| * Licensed to the Apache Software Foundation (ASF) under one or more | ||
| * contributor license agreements. See the NOTICE file distributed with | ||
| * this work for additional information regarding copyright ownership. | ||
| * The ASF licenses this file to You under the Apache License, Version 2.0 | ||
| * (the "License"); you may not use this file except in compliance with | ||
| * the License. You may obtain a copy of the License at | ||
| * | ||
| * http://www.apache.org/licenses/LICENSE-2.0 | ||
| * | ||
| * Unless required by applicable law or agreed to in writing, software | ||
| * distributed under the License is distributed on an "AS IS" BASIS, | ||
| * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
| * See the License for the specific language governing permissions and | ||
| * limitations under the License. | ||
| */ | ||
| package org.apache.kafka.common.errors; | ||
|
|
||
| /** | ||
| * This exception indicates that the size of the telemetry metrics data is too large. | ||
| */ | ||
| public class TelemetryTooLargeException extends ApiException { | ||
|
|
||
| public TelemetryTooLargeException(String message) { | ||
| super(message); | ||
| } | ||
| } | ||
|
|
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,28 @@ | ||
| /* | ||
| * Licensed to the Apache Software Foundation (ASF) under one or more | ||
| * contributor license agreements. See the NOTICE file distributed with | ||
| * this work for additional information regarding copyright ownership. | ||
| * The ASF licenses this file to You under the Apache License, Version 2.0 | ||
| * (the "License"); you may not use this file except in compliance with | ||
| * the License. You may obtain a copy of the License at | ||
| * | ||
| * http://www.apache.org/licenses/LICENSE-2.0 | ||
| * | ||
| * Unless required by applicable law or agreed to in writing, software | ||
| * distributed under the License is distributed on an "AS IS" BASIS, | ||
| * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
| * See the License for the specific language governing permissions and | ||
| * limitations under the License. | ||
| */ | ||
| package org.apache.kafka.common.errors; | ||
|
|
||
| /** | ||
| * This exception indicates that the client sent an invalid or outdated SubscriptionId | ||
| */ | ||
| public class UnknownSubscriptionIdException extends ApiException { | ||
|
|
||
| public UnknownSubscriptionIdException(String message) { | ||
| super(message); | ||
| } | ||
| } | ||
|
|
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -14,19 +14,21 @@ | |
| * See the License for the specific language governing permissions and | ||
| * limitations under the License. | ||
| */ | ||
|
|
||
| package org.apache.kafka.common.requests; | ||
|
|
||
| import org.apache.kafka.common.message.PushTelemetryRequestData; | ||
| import org.apache.kafka.common.message.PushTelemetryResponseData; | ||
| import org.apache.kafka.common.protocol.ApiKeys; | ||
| import org.apache.kafka.common.protocol.ByteBufferAccessor; | ||
| import org.apache.kafka.common.protocol.Errors; | ||
| import org.apache.kafka.common.record.CompressionType; | ||
|
|
||
| import java.nio.ByteBuffer; | ||
|
|
||
| public class PushTelemetryRequest extends AbstractRequest { | ||
|
|
||
| private static final String OTLP_CONTENT_TYPE = "OTLP"; | ||
|
|
||
| public static class Builder extends AbstractRequest.Builder<PushTelemetryRequest> { | ||
|
|
||
| private final PushTelemetryRequestData data; | ||
|
|
@@ -71,6 +73,31 @@ public PushTelemetryRequestData data() { | |
| return data; | ||
| } | ||
|
|
||
| public PushTelemetryResponse createResponse(int throttleTimeMs, Errors errors) { | ||
| PushTelemetryResponseData responseData = new PushTelemetryResponseData(); | ||
| responseData.setErrorCode(errors.code()); | ||
| responseData.setThrottleTimeMs(throttleTimeMs); | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. If error code is not Could we redirect
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Done
It might be naive but sorry I didn't understand as why throttleTimeMs should not be passed in response for other exceptions. Isn't all requests goes through common throttling code where requests might be throttled for sometime based on the throughput?
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The following is my understanding. There are two types of throttling.
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Thanks for the explanation @junrao, this is helpful. I have removed throttleTimeMs for the case when THROTTLING_QUOTA_EXCEEDED exception is thrown in telemetry APIs. Done. |
||
| return new PushTelemetryResponse(responseData); | ||
| } | ||
|
|
||
| public String getMetricsContentType() { | ||
|
apoorvmittal10 marked this conversation as resolved.
Outdated
|
||
| // Future versions of PushTelemetryRequest and GetTelemetrySubscriptionsRequest may include a content-type | ||
| // field to allow for updated OTLP format versions (or additional formats), but this field is currently not | ||
| // included since only one format is specified in the current proposal of the kip-714 | ||
| return OTLP_CONTENT_TYPE; | ||
| } | ||
|
|
||
| public ByteBuffer getMetricsData() { | ||
| CompressionType cType = CompressionType.forId(this.data.compressionType()); | ||
| return (cType == CompressionType.NONE) ? | ||
| ByteBuffer.wrap(this.data.metrics()) : decompressMetricsData(cType, this.data.metrics()); | ||
| } | ||
|
|
||
| private static ByteBuffer decompressMetricsData(CompressionType compressionType, byte[] metrics) { | ||
| // TODO: Add support for decompression of metrics data | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. So, this will be added in a future PR?
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Yes, I have created jira in parent KIP-714 task to address this: https://issues.apache.org/jira/browse/KAFKA-15807. I am planning to get end-to-end metrics flow without compression first. |
||
| return ByteBuffer.wrap(metrics); | ||
| } | ||
|
|
||
| public static PushTelemetryRequest parse(ByteBuffer buffer, short version) { | ||
| return new PushTelemetryRequest(new PushTelemetryRequestData( | ||
| new ByteBufferAccessor(buffer), version), version); | ||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.