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 @@ -23,6 +23,8 @@
import io.netty.handler.ssl.SslHandler;
import io.netty.handler.timeout.IdleStateHandler;
import io.streamnative.pulsar.handlers.kop.stats.StatsLogger;
import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperation;
import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationPurgatory;
import io.streamnative.pulsar.handlers.kop.utils.ssl.SSLUtils;
import java.util.concurrent.TimeUnit;
import lombok.Getter;
Expand All @@ -46,6 +48,8 @@ public class KafkaChannelInitializer extends ChannelInitializer<SocketChannel> {
private final KopBrokerLookupManager kopBrokerLookupManager;

private final AdminManager adminManager;
private DelayedOperationPurgatory<DelayedOperation> producePurgatory;
private DelayedOperationPurgatory<DelayedOperation> fetchPurgatory;
@Getter
private final boolean enableTls;
@Getter
Expand All @@ -60,6 +64,8 @@ public KafkaChannelInitializer(PulsarService pulsarService,
TenantContextManager tenantContextManager,
KopBrokerLookupManager kopBrokerLookupManager,
AdminManager adminManager,
DelayedOperationPurgatory<DelayedOperation> producePurgatory,
DelayedOperationPurgatory<DelayedOperation> fetchPurgatory,
boolean enableTLS,
EndPoint advertisedEndPoint,
StatsLogger statsLogger) {
Expand All @@ -69,6 +75,8 @@ public KafkaChannelInitializer(PulsarService pulsarService,
this.tenantContextManager = tenantContextManager;
this.kopBrokerLookupManager = kopBrokerLookupManager;
this.adminManager = adminManager;
this.producePurgatory = producePurgatory;
this.fetchPurgatory = fetchPurgatory;
this.enableTls = enableTLS;
this.advertisedEndPoint = advertisedEndPoint;
this.statsLogger = statsLogger;
Expand Down Expand Up @@ -100,6 +108,16 @@ protected void initChannel(SocketChannel ch) throws Exception {
public KafkaRequestHandler newCnx() throws Exception {
return new KafkaRequestHandler(pulsarService, kafkaConfig,
tenantContextManager, kopBrokerLookupManager, adminManager,
producePurgatory, fetchPurgatory,
enableTls, advertisedEndPoint, statsLogger);
}

@VisibleForTesting
public KafkaRequestHandler newCnx(final TenantContextManager tenantContextManager,
final StatsLogger statsLogger) throws Exception {
return new KafkaRequestHandler(pulsarService, kafkaConfig,
tenantContextManager, kopBrokerLookupManager, adminManager,
producePurgatory, fetchPurgatory,
enableTls, advertisedEndPoint, statsLogger);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,8 @@
import io.streamnative.pulsar.handlers.kop.utils.ConfigurationUtils;
import io.streamnative.pulsar.handlers.kop.utils.KopTopic;
import io.streamnative.pulsar.handlers.kop.utils.MetadataUtils;
import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperation;
import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationPurgatory;
import io.streamnative.pulsar.handlers.kop.utils.timer.SystemTimer;
import java.net.InetSocketAddress;
import java.util.Map;
Expand Down Expand Up @@ -77,6 +79,8 @@ public class KafkaProtocolHandler implements ProtocolHandler, TenantContextManag
private KopBrokerLookupManager kopBrokerLookupManager;
private AdminManager adminManager = null;
private SystemTopicClient txnTopicClient;
private DelayedOperationPurgatory<DelayedOperation> producePurgatory;
private DelayedOperationPurgatory<DelayedOperation> fetchPurgatory;
@VisibleForTesting
@Getter
private Map<InetSocketAddress, ChannelInitializer<SocketChannel>> channelInitializerMap;
Expand Down Expand Up @@ -547,6 +551,8 @@ private KafkaChannelInitializer newKafkaChannelInitializer(final EndPoint endPoi
this,
kopBrokerLookupManager,
adminManager,
producePurgatory,
fetchPurgatory,
endPoint.isTlsEnabled(),
endPoint,
scopeStatsLogger);
Expand All @@ -558,6 +564,15 @@ public Map<InetSocketAddress, ChannelInitializer<SocketChannel>> newChannelIniti
checkState(kafkaConfig != null);
checkState(brokerService != null);

producePurgatory = DelayedOperationPurgatory.<DelayedOperation>builder()
.purgatoryName("produce")
.timeoutTimer(SystemTimer.builder().executorName("produce").build())
.build();
fetchPurgatory = DelayedOperationPurgatory.<DelayedOperation>builder()
.purgatoryName("fetch")
.timeoutTimer(SystemTimer.builder().executorName("fetch").build())
.build();

try {
ImmutableMap.Builder<InetSocketAddress, ChannelInitializer<SocketChannel>> builder =
ImmutableMap.builder();
Expand All @@ -577,9 +592,21 @@ public Map<InetSocketAddress, ChannelInitializer<SocketChannel>> newChannelIniti
@Override
public void close() {
Optional.ofNullable(LOOKUP_CLIENT_MAP.remove(brokerService.pulsar())).ifPresent(LookupClient::close);
offsetTopicClient.close();
txnTopicClient.close();
adminManager.shutdown();
if (offsetTopicClient != null) {
offsetTopicClient.close();
}
if (txnTopicClient != null) {
txnTopicClient.close();
}
if (adminManager != null) {
adminManager.shutdown();
}
if (producePurgatory != null) {
producePurgatory.shutdown();
}
if (fetchPurgatory != null) {
fetchPurgatory.shutdown();
}
groupCoordinatorsByTenant.values().forEach(GroupCoordinator::shutdown);
kopEventManager.close();
transactionCoordinatorByTenant.values().forEach(TransactionCoordinator::shutdown);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,6 @@
import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperation;
import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationKey;
import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationPurgatory;
import io.streamnative.pulsar.handlers.kop.utils.timer.SystemTimer;
import java.net.InetSocketAddress;
import java.nio.ByteBuffer;
import java.util.ArrayList;
Expand Down Expand Up @@ -219,16 +218,8 @@ public class KafkaRequestHandler extends KafkaCommandDecoder {
// is found.
private final Map<TopicPartition, PendingTopicFutures> pendingTopicFuturesMap = new ConcurrentHashMap<>();
// DelayedOperation for produce and fetch
private final DelayedOperationPurgatory<DelayedOperation> producePurgatory =
DelayedOperationPurgatory.<DelayedOperation>builder()
.purgatoryName("produce")
.timeoutTimer(SystemTimer.builder().executorName("produce").build())
.build();
private final DelayedOperationPurgatory<DelayedOperation> fetchPurgatory =
DelayedOperationPurgatory.<DelayedOperation>builder()
.purgatoryName("fetch")
.timeoutTimer(SystemTimer.builder().executorName("fetch").build())
.build();
private final DelayedOperationPurgatory<DelayedOperation> producePurgatory;
private final DelayedOperationPurgatory<DelayedOperation> fetchPurgatory;

// Flag to manage throttling-publish-buffer by atomically enable/disable read-channel.
private final long maxPendingBytes;
Expand Down Expand Up @@ -285,6 +276,8 @@ public KafkaRequestHandler(PulsarService pulsarService,
TenantContextManager tenantContextManager,
KopBrokerLookupManager kopBrokerLookupManager,
AdminManager adminManager,
DelayedOperationPurgatory<DelayedOperation> producePurgatory,
DelayedOperationPurgatory<DelayedOperation> fetchPurgatory,
Boolean tlsEnabled,
EndPoint advertisedEndPoint,
StatsLogger statsLogger) throws Exception {
Expand All @@ -306,6 +299,8 @@ public KafkaRequestHandler(PulsarService pulsarService,
? new SimpleAclAuthorizer(pulsarService)
: null;
this.adminManager = adminManager;
this.producePurgatory = producePurgatory;
this.fetchPurgatory = fetchPurgatory;
this.tlsEnabled = tlsEnabled;
this.advertisedEndPoint = advertisedEndPoint;
this.topicManager = new KafkaTopicManager(this);
Expand Down Expand Up @@ -355,8 +350,6 @@ protected void close() {
log.info("currentConnectedGroup remove {}", clientHost);
currentConnectedGroup.remove(clientHost);
}
producePurgatory.shutdown();
fetchPurgatory.shutdown();

// update alive channel count stat
RequestStats.ALIVE_CHANNEL_COUNT_INSTANCE.decrementAndGet();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,26 +18,18 @@

import io.netty.channel.Channel;
import io.netty.channel.ChannelHandlerContext;
import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupCoordinator;
import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator;
import io.streamnative.pulsar.handlers.kop.stats.NullStatsLogger;
import java.net.InetSocketAddress;
import java.net.SocketAddress;
import org.apache.pulsar.broker.protocol.ProtocolHandler;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.testng.annotations.AfterMethod;
import org.testng.annotations.BeforeMethod;

/**
* Test for entry publish time.
*/
public class EntryPublishTimeTest extends KopProtocolHandlerTestBase {
private static final Logger log = LoggerFactory.getLogger(EntryPublishTimeTest.class);

KafkaRequestHandler kafkaRequestHandler;
SocketAddress serviceAddress;
private AdminManager adminManager;

public EntryPublishTimeTest(String format) {
super(format);
Expand All @@ -48,32 +40,7 @@ public EntryPublishTimeTest(String format) {
protected void setup() throws Exception {
super.internalSetup();

ProtocolHandler handler = pulsar.getProtocolHandlers().protocol("kafka");
GroupCoordinator groupCoordinator = ((KafkaProtocolHandler) handler)
.getGroupCoordinator(conf.getKafkaMetadataTenant());
TransactionCoordinator transactionCoordinator = ((KafkaProtocolHandler) handler)
.getTransactionCoordinator(conf.getKafkaMetadataTenant());

adminManager = new AdminManager(pulsar.getAdminClient(), conf);
kafkaRequestHandler = new KafkaRequestHandler(
pulsar,
conf,
new TenantContextManager() {
@Override
public GroupCoordinator getGroupCoordinator(String tenant) {
return groupCoordinator;
}

@Override
public TransactionCoordinator getTransactionCoordinator(String tenant) {
return transactionCoordinator;
}
},
((KafkaProtocolHandler) handler).getKopBrokerLookupManager(),
adminManager,
false,
getPlainEndPoint(),
NullStatsLogger.INSTANCE);
kafkaRequestHandler = newRequestHandler();
ChannelHandlerContext mockCtx = mock(ChannelHandlerContext.class);
Channel mockChannel = mock(Channel.class);
doReturn(mockChannel).when(mockCtx).channel();
Expand All @@ -85,7 +52,6 @@ public TransactionCoordinator getTransactionCoordinator(String tenant) {
@AfterMethod
@Override
protected void cleanup() throws Exception {
adminManager.shutdown();
super.internalCleanup();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -30,9 +30,6 @@
import io.netty.channel.Channel;
import io.netty.channel.ChannelHandlerContext;
import io.streamnative.pulsar.handlers.kop.KafkaCommandDecoder.KafkaHeaderAndRequest;
import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupCoordinator;
import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator;
import io.streamnative.pulsar.handlers.kop.stats.NullStatsLogger;
import java.net.InetSocketAddress;
import java.net.SocketAddress;
import java.nio.ByteBuffer;
Expand Down Expand Up @@ -78,7 +75,6 @@
import org.apache.kafka.common.requests.RequestHeader;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.apache.pulsar.broker.protocol.ProtocolHandler;
import org.apache.pulsar.common.policies.data.RetentionPolicies;
import org.testng.annotations.AfterMethod;
import org.testng.annotations.BeforeMethod;
Expand All @@ -93,7 +89,6 @@ public class KafkaApisTest extends KopProtocolHandlerTestBase {

KafkaRequestHandler kafkaRequestHandler;
SocketAddress serviceAddress;
private AdminManager adminManager;

@Override
protected void resetConfig() {
Expand Down Expand Up @@ -123,32 +118,7 @@ protected void setup() throws Exception {

log.info("created namespaces, init handler");

ProtocolHandler handler = pulsar.getProtocolHandlers().protocol("kafka");
GroupCoordinator groupCoordinator = ((KafkaProtocolHandler) handler)
.getGroupCoordinator(conf.getKafkaMetadataTenant());
TransactionCoordinator transactionCoordinator = ((KafkaProtocolHandler) handler)
.getTransactionCoordinator(conf.getKafkaMetadataTenant());

adminManager = new AdminManager(pulsar.getAdminClient(), conf);
kafkaRequestHandler = new KafkaRequestHandler(
pulsar,
(KafkaServiceConfiguration) conf,
new TenantContextManager() {
@Override
public GroupCoordinator getGroupCoordinator(String tenant) {
return groupCoordinator;
}

@Override
public TransactionCoordinator getTransactionCoordinator(String tenant) {
return transactionCoordinator;
}
},
((KafkaProtocolHandler) handler).getKopBrokerLookupManager(),
adminManager,
false,
getPlainEndPoint(),
NullStatsLogger.INSTANCE);
kafkaRequestHandler = newRequestHandler();
ChannelHandlerContext mockCtx = mock(ChannelHandlerContext.class);
Channel mockChannel = mock(Channel.class);
doReturn(mockChannel).when(mockCtx).channel();
Expand All @@ -160,7 +130,6 @@ public TransactionCoordinator getTransactionCoordinator(String tenant) {
@AfterMethod
@Override
protected void cleanup() throws Exception {
adminManager.shutdown();
super.internalCleanup();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,12 +34,9 @@
import io.netty.channel.ChannelHandlerContext;
import io.streamnative.pulsar.handlers.kop.KafkaCommandDecoder.KafkaHeaderAndRequest;
import io.streamnative.pulsar.handlers.kop.KafkaCommandDecoder.KafkaHeaderAndResponse;
import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupCoordinator;
import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupMetadata;
import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupMetadataManager;
import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator;
import io.streamnative.pulsar.handlers.kop.offset.OffsetAndMetadata;
import io.streamnative.pulsar.handlers.kop.stats.NullStatsLogger;
import io.streamnative.pulsar.handlers.kop.utils.TopicNameUtils;
import java.net.InetSocketAddress;
import java.nio.ByteBuffer;
Expand Down Expand Up @@ -104,7 +101,6 @@
import org.apache.kafka.common.serialization.IntegerSerializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.apache.kafka.common.utils.Time;
import org.apache.pulsar.broker.protocol.ProtocolHandler;
import org.apache.pulsar.client.admin.PulsarAdminException;
import org.apache.pulsar.common.allocator.PulsarByteBufAllocator;
import org.apache.pulsar.common.naming.TopicName;
Expand All @@ -123,7 +119,6 @@
public class KafkaRequestHandlerTest extends KopProtocolHandlerTestBase {

private KafkaRequestHandler handler;
private AdminManager adminManager;

@DataProvider(name = "metadataVersions")
public static Object[][] metadataVersions() {
Expand Down Expand Up @@ -153,32 +148,7 @@ protected void setup() throws Exception {

log.info("created namespaces, init handler");

ProtocolHandler handler1 = pulsar.getProtocolHandlers().protocol("kafka");
GroupCoordinator groupCoordinator = ((KafkaProtocolHandler) handler1)
.getGroupCoordinator(conf.getKafkaMetadataTenant());
TransactionCoordinator transactionCoordinator = ((KafkaProtocolHandler) handler1)
.getTransactionCoordinator(conf.getKafkaMetadataTenant());

adminManager = new AdminManager(pulsar.getAdminClient(), conf);
handler = new KafkaRequestHandler(
pulsar,
conf,
new TenantContextManager() {
@Override
public GroupCoordinator getGroupCoordinator(String tenant) {
return groupCoordinator;
}

@Override
public TransactionCoordinator getTransactionCoordinator(String tenant) {
return transactionCoordinator;
}
},
((KafkaProtocolHandler) handler1).getKopBrokerLookupManager(),
adminManager,
false,
getPlainEndPoint(),
NullStatsLogger.INSTANCE);
handler = newRequestHandler();
ChannelHandlerContext mockCtx = mock(ChannelHandlerContext.class);
Channel mockChannel = mock(Channel.class);
doReturn(mockChannel).when(mockCtx).channel();
Expand All @@ -188,7 +158,6 @@ public TransactionCoordinator getTransactionCoordinator(String tenant) {
@AfterClass
@Override
protected void cleanup() throws Exception {
adminManager.shutdown();
super.internalCleanup();
}

Expand Down
Loading