Skip to content
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 @@ -51,7 +51,6 @@
import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.client.api.transaction.Transaction;
import org.apache.pulsar.client.api.transaction.TransactionCoordinatorClientException;
import org.apache.pulsar.client.impl.PulsarClientImpl;
import org.apache.pulsar.client.impl.auth.AuthenticationToken;
import org.apache.pulsar.common.naming.NamespaceName;
import org.apache.pulsar.common.policies.data.AuthAction;
Expand Down Expand Up @@ -208,8 +207,9 @@ public void testEndTxn(String actor, boolean afterUnload) throws Exception {
TransactionCoordinatorID.get(transaction.getTxnID().getMostSigBits()));
}

final Throwable ex = syncGetException((
(PulsarClientImpl) pulsarClientOther).getTcClient().commitAsync(transaction.getTxnID())
final Throwable ex = syncGetException(pulsarClientOther
.getTransactionCoordinatorClient()
.commitAsync(transaction.getTxnID())
);
if (actor.equals("client") || actor.equals("admin")) {
Assert.assertNull(ex);
Expand Down Expand Up @@ -246,11 +246,13 @@ public void testAddPartitionToTxn(String actor, boolean afterUnload) throws Exce
TransactionCoordinatorID.get(transaction.getTxnID().getMostSigBits()));
}

final Throwable ex = syncGetException(((PulsarClientImpl) pulsarClientOther)
.getTcClient().addPublishPartitionToTxnAsync(transaction.getTxnID(), List.of(TOPIC)));
final Throwable ex = syncGetException(pulsarClientOther
.getTransactionCoordinatorClient()
.addPublishPartitionToTxnAsync(transaction.getTxnID(),
List.of(TOPIC)));

final TxnMeta txnMeta = pulsarServiceList.get(0).getTransactionMetadataStoreService()
.getTxnMeta(transaction.getTxnID()).get();
.getTxnMeta(transaction.getTxnID()).get();
if (actor.equals("client") || actor.equals("admin")) {
Assert.assertNull(ex);
Assert.assertEquals(txnMeta.producedPartitions(), List.of(TOPIC));
Expand Down Expand Up @@ -284,8 +286,9 @@ public void testAddSubscriptionToTxn(String actor, boolean afterUnload) throws E
TransactionCoordinatorID.get(transaction.getTxnID().getMostSigBits()));
}

final Throwable ex = syncGetException(((PulsarClientImpl) pulsarClientOther)
.getTcClient().addSubscriptionToTxnAsync(transaction.getTxnID(), TOPIC, "sub"));
final Throwable ex = syncGetException(pulsarClientOther
.getTransactionCoordinatorClient()
.addSubscriptionToTxnAsync(transaction.getTxnID(), TOPIC, "sub"));

final TxnMeta txnMeta = pulsarServiceList.get(0).getTransactionMetadataStoreService()
.getTxnMeta(transaction.getTxnID()).get();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,6 @@
import org.apache.pulsar.client.admin.PulsarAdminException;
import org.apache.pulsar.client.api.PulsarClientException;
import org.apache.pulsar.client.api.transaction.TransactionCoordinatorClient;
import org.apache.pulsar.client.impl.PulsarClientImpl;
import org.apache.pulsar.common.naming.TopicName;
import org.awaitility.Awaitility;
import org.testng.annotations.AfterMethod;
Expand All @@ -45,8 +44,9 @@ public class TransactionBufferCloseTest extends TransactionTestBase {
@BeforeMethod
protected void setup() throws Exception {
setUpBase(1, 16, null, 0);
Awaitility.await().until(() -> ((PulsarClientImpl) pulsarClient)
.getTcClient().getState() == TransactionCoordinatorClient.State.READY);
Awaitility.await().until(() -> pulsarClient
.getTransactionCoordinatorClient()
.getState() == TransactionCoordinatorClient.State.READY);
}

@AfterMethod(alwaysRun = true)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,6 @@
import org.apache.pulsar.client.api.transaction.Transaction;
import org.apache.pulsar.client.api.transaction.TransactionCoordinatorClient;
import org.apache.pulsar.client.impl.MessageIdImpl;
import org.apache.pulsar.client.impl.PulsarClientImpl;
import org.apache.pulsar.common.naming.TopicName;
import org.awaitility.Awaitility;
import org.testng.annotations.AfterMethod;
Expand All @@ -63,8 +62,9 @@ public class TransactionStablePositionTest extends TransactionTestBase {
@BeforeMethod
protected void setup() throws Exception {
setUpBase(1, 16, TOPIC, 0);
Awaitility.await().until(() -> ((PulsarClientImpl) pulsarClient)
.getTcClient().getState() == TransactionCoordinatorClient.State.READY);
Awaitility.await().until(() -> pulsarClient
.getTransactionCoordinatorClient()
.getState() == TransactionCoordinatorClient.State.READY);
}

@AfterMethod(alwaysRun = true)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@
import org.apache.pulsar.client.api.MessageId;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.PulsarClientException;
import org.apache.pulsar.client.api.transaction.TransactionCoordinatorClient;
import org.apache.pulsar.client.api.transaction.TransactionCoordinatorClientException;
import org.apache.pulsar.client.api.transaction.TxnID;
import org.apache.pulsar.client.impl.transaction.TransactionCoordinatorClientImpl;
Expand Down Expand Up @@ -76,7 +77,7 @@ public void testTransactionNewReconnect() throws Exception {

@Test
public void testTransactionAddSubscriptionToTxnAsyncReconnect() throws Exception {
TransactionCoordinatorClientImpl transactionCoordinatorClient = ((PulsarClientImpl) pulsarClient).getTcClient();
TransactionCoordinatorClient transactionCoordinatorClient = pulsarClient.getTransactionCoordinatorClient();
Callable<CompletableFuture<?>> callable = () -> transactionCoordinatorClient
.addSubscriptionToTxnAsync(new TxnID(0, 0), "test", "test");
tryCommandReconnect(callable, callable);
Expand Down Expand Up @@ -108,7 +109,7 @@ public void tryCommandReconnect(Callable<CompletableFuture<?>> callable1, Callab

@Test
public void testTransactionAbortToTxnAsyncReconnect() throws Exception {
TransactionCoordinatorClientImpl transactionCoordinatorClient = ((PulsarClientImpl) pulsarClient).getTcClient();
TransactionCoordinatorClient transactionCoordinatorClient = pulsarClient.getTransactionCoordinatorClient();
Callable<CompletableFuture<?>> callable1 = () -> transactionCoordinatorClient.abortAsync(new TxnID(0,
0));
Callable<CompletableFuture<?>> callable2 = () -> transactionCoordinatorClient.abortAsync(new TxnID(0,
Expand All @@ -118,7 +119,7 @@ public void testTransactionAbortToTxnAsyncReconnect() throws Exception {

@Test
public void testTransactionCommitToTxnAsyncReconnect() throws Exception {
TransactionCoordinatorClientImpl transactionCoordinatorClient = ((PulsarClientImpl) pulsarClient).getTcClient();
TransactionCoordinatorClient transactionCoordinatorClient = pulsarClient.getTransactionCoordinatorClient();
Callable<CompletableFuture<?>> callable1 = () -> transactionCoordinatorClient.commitAsync(new TxnID(0,
0));
Callable<CompletableFuture<?>> callable2 = () -> transactionCoordinatorClient.commitAsync(new TxnID(0,
Expand All @@ -128,7 +129,7 @@ public void testTransactionCommitToTxnAsyncReconnect() throws Exception {

@Test
public void testTransactionAddPublishPartitionToTxnReconnect() throws Exception {
TransactionCoordinatorClientImpl transactionCoordinatorClient = ((PulsarClientImpl) pulsarClient).getTcClient();
TransactionCoordinatorClient transactionCoordinatorClient = pulsarClient.getTransactionCoordinatorClient();
Callable<CompletableFuture<?>> callable =
() -> transactionCoordinatorClient.addPublishPartitionToTxnAsync(new TxnID(0, 0),
Collections.singletonList("test"));
Expand All @@ -137,7 +138,8 @@ public void testTransactionAddPublishPartitionToTxnReconnect() throws Exception

@Test
public void testPulsarClientCloseThenCloseTcClient() throws Exception {
TransactionCoordinatorClientImpl transactionCoordinatorClient = ((PulsarClientImpl) pulsarClient).getTcClient();
TransactionCoordinatorClientImpl transactionCoordinatorClient =
(TransactionCoordinatorClientImpl) pulsarClient.getTransactionCoordinatorClient();
java.util.Collection<TransactionMetaStoreHandler> handlers = transactionCoordinatorClient.getHandlers();

for (TransactionMetaStoreHandler handler : handlers) {
Expand All @@ -162,9 +164,9 @@ public void testPulsarClientCloseThenCloseTcClient() throws Exception {
}

@Test
public void testHandlerStateChangeToReady() throws Exception {
public void testHandlerStateChangeToReady() {
TransactionCoordinatorClientImpl transactionCoordinatorClient =
((PulsarClientImpl) pulsarClient).getTcClient();
(TransactionCoordinatorClientImpl) pulsarClient.getTransactionCoordinatorClient();
TransactionMetaStoreHandler transactionMetaStoreHandler =
transactionCoordinatorClient.getHandlers().iterator().next();
Assert.assertEquals(transactionMetaStoreHandler.getConnectHandleState(), HandlerState.State.Ready);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -970,7 +970,7 @@ public void txnMetadataHandlerRecoverTest() throws Exception {
.enableTransaction(true)
.build();

TransactionCoordinatorClient tcClient = recoverPulsarClient.getTcClient();
TransactionCoordinatorClient tcClient = recoverPulsarClient.getTransactionCoordinatorClient();
for (TxnID txnID : txnIDList) {
tcClient.commit(txnID);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
import java.util.List;
import java.util.concurrent.CompletableFuture;
import org.apache.pulsar.client.api.transaction.TransactionBuilder;
import org.apache.pulsar.client.api.transaction.TransactionCoordinatorClient;
import org.apache.pulsar.client.internal.DefaultImplementation;
import org.apache.pulsar.common.classification.InterfaceAudience;
import org.apache.pulsar.common.classification.InterfaceStability;
Expand Down Expand Up @@ -400,4 +401,11 @@ default CompletableFuture<List<String>> getPartitionsForTopic(String topic) {
* @since 2.7.0
*/
TransactionBuilder newTransaction();

/**
* Returns the {@link TransactionCoordinatorClient} used by this client.
* This is a lower-level API, meant mainly for usage in streaming engines.
* @return the Transaction Coordinator Client associated with this Pulsar Client
*/
TransactionCoordinatorClient getTransactionCoordinatorClient();
}
Original file line number Diff line number Diff line change
Expand Up @@ -22,13 +22,11 @@
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
import org.apache.pulsar.common.classification.InterfaceAudience;
import org.apache.pulsar.common.classification.InterfaceStability;

/**
* Transaction coordinator client.
*/
@InterfaceAudience.Private
@InterfaceStability.Evolving
public interface TransactionCoordinatorClient extends Closeable {

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -77,12 +77,14 @@
import org.apache.pulsar.client.api.schema.KeyValueSchema;
import org.apache.pulsar.client.api.schema.SchemaInfoProvider;
import org.apache.pulsar.client.api.transaction.TransactionBuilder;
import org.apache.pulsar.client.api.transaction.TransactionCoordinatorClient;
import org.apache.pulsar.client.api.v5.auth.BinaryAuthDataProvider;
import org.apache.pulsar.client.api.v5.internal.ClientAuthenticationServices;
import org.apache.pulsar.client.api.v5.internal.ClientAuthenticationServicesAware;
import org.apache.pulsar.client.impl.auth.oauth2.AuthenticationOAuth2;
import org.apache.pulsar.client.impl.auth.v5.DefaultClientAuthenticationServices;
import org.apache.pulsar.client.impl.auth.v5.FrameworkHttpClientFactory;
import org.apache.pulsar.client.impl.auth.v5.LegacyV4AuthenticationAdapter;
import org.apache.pulsar.client.impl.auth.v5.V5AuthenticationLoader;
import org.apache.pulsar.client.impl.auth.v5.V5BinaryAuthenticationDriver;
import org.apache.pulsar.client.impl.conf.ClientConfigurationData;
Expand Down Expand Up @@ -205,7 +207,7 @@

private final LoadingCache<String, SchemaInfoProvider> schemaProviderLoadingCache =
CacheBuilder.newBuilder().maximumSize(100000)
.expireAfterAccess(30, TimeUnit.MINUTES)

Check warning on line 210 in pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Protobuf v3

[deprecation] expireAfterAccess(long,TimeUnit) in CacheBuilder has been deprecated

Check warning on line 210 in pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java

View workflow job for this annotation

GitHub Actions / Build Pulsar on MacOS

[deprecation] expireAfterAccess(long,TimeUnit) in CacheBuilder has been deprecated
.build(new CacheLoader<String, SchemaInfoProvider>() {

@Override
Expand Down Expand Up @@ -2104,4 +2106,9 @@
NameResolver<InetAddress> getNameResolver() {
return DnsResolverUtil.adaptToNameResolver(addressResolver);
}

@Override
public TransactionCoordinatorClient getTransactionCoordinatorClient() {
return tcClient;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
import org.apache.pulsar.client.api.MessageId;
import org.apache.pulsar.client.api.PulsarClientException;
import org.apache.pulsar.client.api.transaction.Transaction;
import org.apache.pulsar.client.api.transaction.TransactionCoordinatorClient;
import org.apache.pulsar.client.api.transaction.TransactionCoordinatorClientException.InvalidTxnStatusException;
import org.apache.pulsar.client.api.transaction.TransactionCoordinatorClientException.TransactionNotFoundException;
import org.apache.pulsar.client.api.transaction.TxnID;
Expand Down Expand Up @@ -62,7 +63,7 @@ public class TransactionImpl implements Transaction , TimerTask {

private final Map<String, CompletableFuture<Void>> registerPartitionMap;
private final Map<Pair<String, String>, CompletableFuture<Void>> registerSubscriptionMap;
private final TransactionCoordinatorClientImpl tcClient;
private final TransactionCoordinatorClient tcClient;

private CompletableFuture<Void> opFuture;

Expand Down Expand Up @@ -96,7 +97,7 @@ public void run(Timeout timeout) throws Exception {

this.registerPartitionMap = new ConcurrentHashMap<>();
this.registerSubscriptionMap = new ConcurrentHashMap<>();
this.tcClient = client.getTcClient();
this.tcClient = client.getTransactionCoordinatorClient();
this.opFuture = CompletableFuture.completedFuture(null);
this.timeout = client.getTimer().newTimeout(this, transactionTimeoutMs, TimeUnit.MILLISECONDS);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,6 @@
import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.client.api.SubscriptionInitialPosition;
import org.apache.pulsar.client.api.transaction.Transaction;
import org.apache.pulsar.client.impl.PulsarClientImpl;
import org.apache.pulsar.client.impl.transaction.TransactionCoordinatorClientImpl;
import org.apache.pulsar.tests.integration.containers.BrokerContainer;
import org.testng.annotations.Test;
Expand Down Expand Up @@ -138,7 +137,8 @@ public void v4SdkTransactionStillUsesLegacyCoordinator() throws Exception {

// Routing assertion: a v4 SDK client must NOT use metadata-store discovery even though the
// cluster has the v5 coordinator enabled — it stays on the legacy assign-topic coordinator.
TransactionCoordinatorClientImpl tcClient = ((PulsarClientImpl) client).getTcClient();
TransactionCoordinatorClientImpl tcClient =
(TransactionCoordinatorClientImpl) client.getTransactionCoordinatorClient();
assertFalse(tcClient.isUsingMetadataDiscovery(),
"v4 SDK client must use the legacy coordinator, not metadata-store discovery");

Expand Down
Loading