Skip to content
This repository was archived by the owner on Jan 24, 2024. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,6 @@
import lombok.extern.slf4j.Slf4j;
import org.apache.bookkeeper.mledger.Position;
import org.apache.bookkeeper.mledger.impl.PositionImpl;
import org.apache.kafka.clients.admin.NewPartitions;
import org.apache.kafka.common.Node;
import org.apache.kafka.common.config.ConfigResource;
import org.apache.kafka.common.errors.InvalidPartitionsException;
Expand All @@ -49,6 +48,7 @@
import org.apache.kafka.common.errors.UnknownTopicOrPartitionException;
import org.apache.kafka.common.protocol.Errors;
import org.apache.kafka.common.requests.ApiError;
import org.apache.kafka.common.requests.CreatePartitionsRequest;
import org.apache.kafka.common.requests.CreateTopicsRequest;
import org.apache.kafka.common.requests.DescribeConfigsResponse;
import org.apache.kafka.common.requests.MetadataResponse;
Expand Down Expand Up @@ -299,9 +299,10 @@ public void truncateTopic(String topicToDelete,

}

CompletableFuture<Map<String, ApiError>> createPartitionsAsync(Map<String, NewPartitions> createInfo,
int timeoutMs,
String namespacePrefix) {
CompletableFuture<Map<String, ApiError>> createPartitionsAsync(
Map<String, CreatePartitionsRequest.PartitionDetails> createInfo,
int timeoutMs,
String namespacePrefix) {
final Map<String, CompletableFuture<ApiError>> futureMap = new ConcurrentHashMap<>();
final AtomicInteger numTopics = new AtomicInteger(createInfo.size());
final CompletableFuture<Map<String, ApiError>> resultFuture = new CompletableFuture<>();
Expand Down Expand Up @@ -336,12 +337,12 @@ CompletableFuture<Map<String, ApiError>> createPartitionsAsync(Map<String, NewPa
if (numTopics.decrementAndGet() == 0) {
complete.run();
}
} else if (newPartitions.assignments() != null
&& !newPartitions.assignments().isEmpty()) {
} else if (newPartitions.newAssignments() != null
&& !newPartitions.newAssignments().isEmpty()) {
errorFuture.complete(ApiError.fromThrowable(
new InvalidRequestException(
"Kop server currently doesn't support manual assignment replica sets '"
+ newPartitions.assignments() + "' the number of partitions must be specified ")
+ newPartitions.newAssignments() + "' the number of partitions must be specified ")
));

if (numTopics.decrementAndGet() == 0) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@

import static com.google.common.base.Preconditions.checkArgument;
import static org.apache.kafka.common.protocol.ApiKeys.API_VERSIONS;
import static org.apache.kafka.common.protocol.ApiKeys.LIST_OFFSETS;

import io.netty.buffer.ByteBuf;
import io.netty.channel.Channel;
Expand Down Expand Up @@ -44,6 +45,7 @@
import org.apache.kafka.common.requests.AbstractResponse;
import org.apache.kafka.common.requests.ApiVersionsRequest;
import org.apache.kafka.common.requests.KopResponseUtils;
import org.apache.kafka.common.requests.ListOffsetRequestV0;
import org.apache.kafka.common.requests.RequestHeader;
import org.apache.kafka.common.requests.ResponseCallbackWrapper;
import org.apache.kafka.common.requests.ResponseHeader;
Expand Down Expand Up @@ -144,12 +146,24 @@ protected KafkaHeaderAndRequest byteBufToRequest(ByteBuf msg,
} else {
ApiKeys apiKey = header.apiKey();
short apiVersion = header.apiVersion();
if (apiKey.equals(LIST_OFFSETS) && apiVersion == 0) {
ListOffsetRequestV0 body = ListOffsetRequestV0.parse(nio, apiVersion);
return new KafkaHeaderAndRequest(header, body, msg, remoteAddress);
}
Struct struct = apiKey.parseRequest(apiVersion, nio);
AbstractRequest body = AbstractRequest.parseRequest(apiKey, apiVersion, struct);
return new KafkaHeaderAndRequest(header, body, msg, remoteAddress);
}
}

protected ListOffsetRequestV0 byteBufToListOffsetRequestV0(ByteBuf buf) {
checkArgument(buf.readableBytes() > 0);
ByteBuffer nio = buf.nioBuffer();
RequestHeader header = RequestHeader.parse(nio);
short apiVersion = header.apiVersion();
return ListOffsetRequestV0.parse(nio, apiVersion);
}

protected static ByteBuf responseToByteBuf(AbstractResponse response, KafkaHeaderAndRequest request) {
try (KafkaHeaderAndResponse kafkaHeaderAndResponse =
KafkaHeaderAndResponse.responseForRequest(request, response)) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,6 @@
import org.apache.commons.collections4.ListUtils;
import org.apache.commons.lang3.NotImplementedException;
import org.apache.commons.lang3.tuple.Pair;
import org.apache.kafka.clients.admin.NewPartitions;
import org.apache.kafka.common.Node;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.acl.AclOperation;
Expand Down Expand Up @@ -143,6 +142,7 @@
import org.apache.kafka.common.requests.LeaveGroupRequest;
import org.apache.kafka.common.requests.ListGroupsRequest;
import org.apache.kafka.common.requests.ListOffsetRequest;
import org.apache.kafka.common.requests.ListOffsetRequestV0;
import org.apache.kafka.common.requests.ListOffsetResponse;
import org.apache.kafka.common.requests.MetadataRequest;
import org.apache.kafka.common.requests.MetadataResponse.PartitionMetadata;
Expand Down Expand Up @@ -1285,7 +1285,7 @@ private void handleListOffsetRequestV1AndAbove(KafkaHeaderAndRequest listOffset,
completeOne.run();
return;
}
responseData.put(topic, fetchOffset(fullPartitionName, times));
responseData.put(topic, fetchOffset(fullPartitionName, times.timestamp));
completeOne.run();
}
);
Expand All @@ -1298,7 +1298,7 @@ private void handleListOffsetRequestV1AndAbove(KafkaHeaderAndRequest listOffset,
// https://cfchou.github.io/blog/2015/04/23/a-closer-look-at-kafka-offsetrequest/ through web.archive.org
private void handleListOffsetRequestV0(KafkaHeaderAndRequest listOffset,
CompletableFuture<AbstractResponse> resultFuture) {
ListOffsetRequest request = (ListOffsetRequest) listOffset.getRequest();
ListOffsetRequestV0 request = (ListOffsetRequestV0) listOffset.getRequest();

Map<TopicPartition, CompletableFuture<Pair<Errors, Long>>> responseData =
Maps.newConcurrentMap();
Expand Down Expand Up @@ -1359,7 +1359,8 @@ private void handleListOffsetRequestV0(KafkaHeaderAndRequest listOffset,
@Override
protected void handleListOffsetRequest(KafkaHeaderAndRequest listOffset,
CompletableFuture<AbstractResponse> resultFuture) {
checkArgument(listOffset.getRequest() instanceof ListOffsetRequest);
checkArgument(listOffset.getRequest() instanceof ListOffsetRequest
|| listOffset.getRequest() instanceof ListOffsetRequestV0);
// the only difference between v0 and v1 is the `max_num_offsets => INT32`
// v0 is required because it is used by librdkafka
if (listOffset.getHeader().apiVersion() == 0) {
Expand Down Expand Up @@ -2451,7 +2452,7 @@ protected void handleCreatePartitions(KafkaHeaderAndRequest createPartitions,
CreatePartitionsRequest request = (CreatePartitionsRequest) createPartitions.getRequest();

final Map<String, ApiError> result = Maps.newConcurrentMap();
final Map<String, NewPartitions> validTopics = Maps.newHashMap();
final Map<String, CreatePartitionsRequest.PartitionDetails> validTopics = Maps.newHashMap();
final Set<String> duplicateTopics = request.duplicates();

KafkaRequestUtils.forEachCreatePartitionsRequest(request, (topic, newPartition) -> {
Expand All @@ -2472,7 +2473,7 @@ protected void handleCreatePartitions(KafkaHeaderAndRequest createPartitions,

String namespacePrefix = currentNamespacePrefix();
final AtomicInteger validTopicsCount = new AtomicInteger(validTopics.size());
final Map<String, NewPartitions> authorizedTopics = Maps.newConcurrentMap();
final Map<String, CreatePartitionsRequest.PartitionDetails> authorizedTopics = Maps.newConcurrentMap();
Runnable createPartitionsAsync = () -> {
if (authorizedTopics.isEmpty()) {
resultFuture.complete(KafkaResponseUtils.newCreatePartitions(result));
Expand All @@ -2486,7 +2487,7 @@ protected void handleCreatePartitions(KafkaHeaderAndRequest createPartitions,
});
};

BiConsumer<String, NewPartitions> completeOneTopic = (topic, newPartitions) -> {
BiConsumer<String, CreatePartitionsRequest.PartitionDetails> completeOneTopic = (topic, newPartitions) -> {
authorizedTopics.put(topic, newPartitions);
if (validTopicsCount.decrementAndGet() == 0) {
createPartitionsAsync.run();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.function.Consumer;
import javax.security.sasl.SaslException;
import lombok.extern.slf4j.Slf4j;
import org.apache.bookkeeper.util.collections.ConcurrentLongHashMap;
import org.apache.kafka.common.protocol.ApiKeys;
Expand Down Expand Up @@ -221,7 +222,7 @@ private CompletableFuture<ChannelHandlerContext> authenticateInternal(ChannelHan
saslAuthBytes = usernamePassword.getBytes(UTF_8);
break;
case OAuthBearerLoginModule.OAUTHBEARER_MECHANISM:
saslAuthBytes = new OAuthBearerClientInitialResponse(commandData).toBytes();
saslAuthBytes = new OAuthBearerClientInitialResponse(commandData, null, null).toBytes();
break;
default:
log.error("No corresponding mechanism to {}", authentication.getClass().getName());
Expand Down Expand Up @@ -252,7 +253,7 @@ private CompletableFuture<ChannelHandlerContext> authenticateInternal(ChannelHan
result.completeExceptionally(saslResponse.error().exception());
}
});
} catch (PulsarClientException ex) {
} catch (PulsarClientException | SaslException ex) {
log.error("Transaction marker channel handler authentication failed.", ex);
result.completeExceptionally(ex);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,37 +18,38 @@
import java.util.function.BiConsumer;
import java.util.function.Consumer;
import java.util.function.Function;
import org.apache.kafka.clients.admin.NewPartitions;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.requests.CreatePartitionsRequest;
import org.apache.kafka.common.requests.ListOffsetRequest;
import org.apache.kafka.common.requests.ListOffsetRequestV0;
import org.apache.kafka.common.requests.OffsetCommitRequest;
import org.apache.kafka.common.requests.TxnOffsetCommitRequest;

public class KafkaRequestUtils {

public static void forEachCreatePartitionsRequest(CreatePartitionsRequest request,
BiConsumer<String, NewPartitions> consumer) {
public static void forEachCreatePartitionsRequest(
CreatePartitionsRequest request,
BiConsumer<String, CreatePartitionsRequest.PartitionDetails> consumer) {
request.newPartitions().forEach(consumer);
}

public static void forEachListOffsetRequest(ListOffsetRequest request,
BiConsumer<TopicPartition, Long> consumer) {
BiConsumer<TopicPartition, ListOffsetRequest.PartitionData> consumer) {
request.partitionTimestamps().forEach(consumer);
}

public static String getMetadata(TxnOffsetCommitRequest.CommittedOffset committedOffset) {
return Optional.ofNullable(committedOffset.metadata()).orElse(OffsetAndMetadata.NoMetadata);
return Optional.ofNullable(committedOffset.metadata).orElse(OffsetAndMetadata.NoMetadata);
}

public static long getOffset(TxnOffsetCommitRequest.CommittedOffset committedOffset) {
return committedOffset.offset();
return committedOffset.offset;
}

public static class LegacyUtils {

public static void forEachListOffsetRequest(
ListOffsetRequest request,
ListOffsetRequestV0 request,
Function<TopicPartition, Function<Long, Consumer<Integer>>> function) {
request.offsetData().forEach((topicPartition, partitionData) -> {
function.apply(topicPartition).apply(partitionData.timestamp).accept(partitionData.maxNumOffsets);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -147,21 +147,24 @@ public static ListGroupsResponse newListGroups(Errors errors,

public static ListOffsetResponse newListOffset(
Map<TopicPartition,
Pair<Errors, Long>> partitionToOffset,
Pair<Errors, Long>> partitionToOffset,
boolean legacy) {
if (legacy) {
return new ListOffsetResponse(CoreUtils.mapValue(partitionToOffset,
pair -> new ListOffsetResponse.PartitionData(
pair.getLeft(),
Optional.ofNullable(pair.getRight()).map(Collections::singletonList)
.orElse(Collections.emptyList()))
));
));
} else {
return new ListOffsetResponse(CoreUtils.mapValue(partitionToOffset,
pair -> new ListOffsetResponse.PartitionData(
pair.getLeft(), // error
0L, // timestamp
Optional.ofNullable(pair.getRight()).orElse(0L) // offset
Optional.ofNullable(
pair.getRight() != null ? pair.getRight().intValue() : null)
.orElse(0) // offset
, Optional.empty()
)
));
}
Expand All @@ -179,6 +182,7 @@ public static MetadataResponse.PartitionMetadata newMetadataPartition(int partit
return new MetadataResponse.PartitionMetadata(Errors.NONE,
partition,
node, // leader
Optional.empty(), // leaderEpoch is unknown in Pulsar
Collections.singletonList(node), // replicas
Collections.singletonList(node), // isr
Collections.emptyList() // offline replicas
Expand All @@ -190,6 +194,7 @@ public static MetadataResponse.PartitionMetadata newMetadataPartition(Errors err
return new MetadataResponse.PartitionMetadata(errors,
partition,
Node.noNode(), // leader
Optional.empty(), // leaderEpoch is unknown in Pulsar
Collections.singletonList(Node.noNode()), // replicas
Collections.singletonList(Node.noNode()), // isr
Collections.emptyList() // offline replicas
Expand All @@ -203,12 +208,14 @@ public static OffsetCommitResponse newOffsetCommit(Map<TopicPartition, Errors> r
public static OffsetFetchResponse.PartitionData newOffsetFetchPartition(long offset,
String metadata) {
return new OffsetFetchResponse.PartitionData(offset,
Optional.empty(), // leaderEpoch is unknown in Pulsar
metadata,
Errors.NONE);
}

public static OffsetFetchResponse.PartitionData newOffsetFetchPartition() {
return new OffsetFetchResponse.PartitionData(OffsetFetchResponse.INVALID_OFFSET,
Optional.empty(), // leaderEpoch is unknown in Pulsar
"", // metadata
Errors.NONE
);
Expand Down
Loading