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 @@ -131,7 +131,8 @@ public static TransactionCoordinator of(String tenant,
return new TransactionCoordinator(
transactionConfig,
new TransactionMarkerChannelManager(tenant, kafkaConfig, transactionStateManager,
kopBrokerLookupManager, false, namespacePrefixForUserTopics),
kopBrokerLookupManager, false, namespacePrefixForUserTopics,
scheduler),
scheduler,
new ProducerIdManagerImpl(transactionConfig.getBrokerId(), metadataStore),
transactionStateManager,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,9 @@
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.BiConsumer;
import lombok.AllArgsConstructor;
Expand Down Expand Up @@ -82,6 +85,8 @@ public class TransactionMarkerChannelManager {
private BlockingQueue<PendingCompleteTxn> txnLogAppendRetryQueue = new LinkedBlockingQueue<>();
private volatile boolean closed;
private final String namespacePrefixForUserTopics;
private final ScheduledExecutorService scheduler;
private ScheduledFuture<?> drainQueuedTransactionMarkersHandle;

@AllArgsConstructor
@ToString
Expand Down Expand Up @@ -141,13 +146,15 @@ public TransactionMarkerChannelManager(String tenant,
TransactionStateManager txnStateManager,
KopBrokerLookupManager kopBrokerLookupManager,
boolean enableTls,
String namespacePrefixForUserTopics) throws Exception {
String namespacePrefixForUserTopics,
ScheduledExecutorService scheduler) throws Exception {
this.tenant = tenant;
this.kafkaConfig = kafkaConfig;
this.namespacePrefixForUserTopics = namespacePrefixForUserTopics;
this.txnStateManager = txnStateManager;
this.kopBrokerLookupManager = kopBrokerLookupManager;
this.enableTls = enableTls;
this.scheduler = scheduler;
if (this.enableTls) {
sslContextFactory = SSLUtils.createSslContextFactory(kafkaConfig);
sslEndPoint = EndPoint.getSslEndPoint(kafkaConfig.getKafkaListeners());
Expand All @@ -171,25 +178,13 @@ public TransactionMarkerChannelManager(String tenant,
bootstrap.group(eventLoopGroup);
bootstrap.channel(NioSocketChannel.class);
bootstrap.handler(new TransactionMarkerChannelInitializer(kafkaConfig, enableTls, this));

Thread thread = new Thread(() -> {
while (!closed) {
drainQueuedTransactionMarkers();
try {
Thread.sleep(1);
} catch (InterruptedException e) {
log.info("ignore {}", e);
}
}
}, "kop-transaction-channel-manager-" + namespacePrefixForUserTopics);
thread.setDaemon(true);
thread.start();
}

public CompletableFuture<TransactionMarkerChannelHandler> getChannel(InetSocketAddress socketAddress) {
if (closed) {
return FutureUtil.failedFuture(new Exception("This TransactionMarkerChannelManager is closed"));
}
ensureDrainQueuedTransactionMarkersActivity();
return handlerMap.computeIfAbsent(socketAddress, address -> {
CompletableFuture<TransactionMarkerChannelHandler> handlerFuture = new CompletableFuture<>();
ChannelFutures.toCompletableFuture(bootstrap.connect(socketAddress))
Expand Down Expand Up @@ -225,6 +220,7 @@ public void addTxnMarkersToSend(Integer coordinatorEpoch,
TransactionMetadata txnMetadata,
TransactionMetadata.TxnTransitMetadata newMetadata,
String namespacePrefix) {
ensureDrainQueuedTransactionMarkersActivity();
String transactionalId = txnMetadata.getTransactionalId();
PendingCompleteTxn pendingCompleteTxn = new PendingCompleteTxn(
transactionalId,
Expand All @@ -249,6 +245,7 @@ private boolean hasPendingMarkersToWrite(TransactionMetadata txnMetadata) {
}

public void maybeWriteTxnCompletion(String transactionalId) {
ensureDrainQueuedTransactionMarkersActivity();
PendingCompleteTxn pendingCompleteTxn = transactionsWithPendingMarkers.get(transactionalId);
if (!hasPendingMarkersToWrite(pendingCompleteTxn.txnMetadata)
&& transactionsWithPendingMarkers.remove(transactionalId, pendingCompleteTxn)) {
Expand All @@ -263,6 +260,7 @@ public void addTxnMarkersToBrokerQueue(String transactionalId,
Integer coordinatorEpoch,
Set<TopicPartition> topicPartitions,
String namespacePrefixForUserTopics) {
ensureDrainQueuedTransactionMarkersActivity();
Integer txnTopicPartition = txnStateManager.partitionFor(transactionalId);

Map<InetSocketAddress, List<TopicPartition>> addressAndPartitionMap = new ConcurrentHashMap<>();
Expand Down Expand Up @@ -404,6 +402,7 @@ public void fail(Errors errors) {
}

public void removeMarkersForTxnTopicPartition(Integer txnTopicPartitionId) {
ensureDrainQueuedTransactionMarkersActivity();
BlockingQueue<TxnIdAndMarkerEntry> unknownBrokerMarkerEntries =
markersQueueForUnknownBroker.removeMarkersForTxnTopicPartition(txnTopicPartitionId);
if (unknownBrokerMarkerEntries != null) {
Expand Down Expand Up @@ -475,8 +474,24 @@ private void drainQueuedTransactionMarkers() {
}
}

private synchronized void ensureDrainQueuedTransactionMarkersActivity() {
if (drainQueuedTransactionMarkersHandle != null || closed) {
return;
}
drainQueuedTransactionMarkersHandle = scheduler.scheduleWithFixedDelay(() -> {
drainQueuedTransactionMarkers();
}, 100, 100, TimeUnit.MILLISECONDS);
Comment thread
BewareMyPower marked this conversation as resolved.
}

private synchronized void stopDrainQueuedTransactionMarkersHandleActivity() {
if (drainQueuedTransactionMarkersHandle != null) {
drainQueuedTransactionMarkersHandle.cancel(false);
}
}

public void close() {
this.closed = true;
stopDrainQueuedTransactionMarkersHandleActivity();
handlerMap.forEach((address, handler) -> {
try {
final TransactionMarkerChannelHandler transactionMarkerChannelHandler = handler.get();
Expand Down