From c7dff9f7535dc075a289840e7693167ed1036d02 Mon Sep 17 00:00:00 2001 From: Madhavan Narayanan Date: Tue, 11 Jan 2022 15:10:50 +0530 Subject: [PATCH 1/2] Interceptors for transaction begin and end events --- .../broker/intercept/BrokerInterceptor.java | 17 +++++++++++++++ .../BrokerInterceptorWithClassLoader.java | 10 +++++++++ .../broker/intercept/BrokerInterceptors.java | 21 +++++++++++++++++++ .../pulsar/broker/service/ServerCnx.java | 9 ++++++++ .../service/persistent/PersistentTopic.java | 1 + 5 files changed, 58 insertions(+) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptor.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptor.java index 1e440b883da28..ee2cc6185e0b9 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptor.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptor.java @@ -123,6 +123,23 @@ default void messageAcked(ServerCnx cnx, Consumer consumer, CommandAck ackCmd) { } + /** + * Intercept when a transaction begins. + * + * @param tcId Transaction Coordinator Id + * @param txnID Transaction ID + */ + default void beginTxn(long tcId, String txnID) { + } + + /** + * Intercept when a transaction ends. + * + * @param txnID Transaction ID + * @param txnAction Transaction Action + */ + default void endTxn(String txnID, long txnAction) { + } /** * Called by the broker while new command incoming. */ diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java index 6725f67bc0180..d0434cc183379 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java @@ -92,6 +92,16 @@ public void messageAcked(ServerCnx cnx, Consumer consumer, this.interceptor.messageAcked(cnx, consumer, ackCmd); } + @Override + public void beginTxn(long tcId, String txnID) { + this.interceptor.beginTxn(tcId, txnID); + } + + @Override + public void endTxn(String txnID, long txnAction) { + this.interceptor.endTxn(txnID, txnAction); + } + @Override public void onConnectionCreated(ServerCnx cnx) { this.interceptor.onConnectionCreated(cnx); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptors.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptors.java index 878f2cb5320a2..75211276f8e89 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptors.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptors.java @@ -166,6 +166,27 @@ public void messageAcked(ServerCnx cnx, Consumer consumer, } } + @Override + public void beginTxn(long tcId, String txnID) { + if (interceptors == null || interceptors.isEmpty()) { + return; + } + for (BrokerInterceptorWithClassLoader value : interceptors.values()) { + value.beginTxn(tcId, txnID); + } + } + + @Override + public void endTxn(String txnID, long txnAction) { + if (interceptors == null || interceptors.isEmpty()) { + return; + } + for (BrokerInterceptorWithClassLoader value : interceptors.values()) { + value.endTxn(txnID, txnAction); + } + } + + @Override public void onConnectionCreated(ServerCnx cnx) { if (interceptors == null || interceptors.isEmpty()) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index 5b4d1f57d7a8b..91b90ea9a61ab 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -2068,6 +2068,9 @@ protected void handleNewTxn(CommandNewTxn command) { } ctx.writeAndFlush(Commands.newTxnResponse(requestId, txnID.getLeastSigBits(), txnID.getMostSigBits())); + if (getBrokerService().getInterceptor() != null) { + getBrokerService().getInterceptor().beginTxn(command.getTcId(), txnID.toString()); + } } else { ex = handleTxnException(ex, BaseCommand.Type.NEW_TXN.name(), requestId); @@ -2135,12 +2138,18 @@ protected void handleEndTxn(CommandEndTxn command) { if (ex == null) { ctx.writeAndFlush(Commands.newEndTxnResponse(requestId, txnID.getLeastSigBits(), txnID.getMostSigBits())); + if (getBrokerService().getInterceptor() != null) { + getBrokerService().getInterceptor().endTxn(txnID.toString(), txnAction); + } } else { ex = handleTxnException(ex, BaseCommand.Type.END_TXN.name(), requestId); ctx.writeAndFlush(Commands.newEndTxnResponse(requestId, txnID.getMostSigBits(), BrokerServiceException.getClientErrorCode(ex), ex.getMessage())); transactionMetadataStoreService.handleOpFail(ex, tcId); + if (getBrokerService().getInterceptor() != null) { + getBrokerService().getInterceptor().endTxn(txnID.toString(), TxnAction.ABORT_VALUE); + } } }); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 5c6f99ac480ba..8e47bc67316e2 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -2997,6 +2997,7 @@ public void publishTxnMessage(TxnID txnID, ByteBuf headersAndPayload, PublishCon // Message has been successfully persisted messageDeduplication.recordMessagePersisted(publishContext, (PositionImpl) position); + publishContext.setProperty("txn_id", txnID.toString()); publishContext.completed(null, ((PositionImpl) position).getLedgerId(), ((PositionImpl) position).getEntryId()); From 50243b758552c3ed1ae8f56dbffcaae21838a7a9 Mon Sep 17 00:00:00 2001 From: Madhavan Narayanan Date: Thu, 17 Mar 2022 12:56:17 +0530 Subject: [PATCH 2/2] pulsar-broker: Changes to address review comments in PR#14613 --- .../broker/intercept/BrokerInterceptor.java | 4 +- .../BrokerInterceptorWithClassLoader.java | 8 ++-- .../broker/intercept/BrokerInterceptors.java | 8 ++-- .../broker/service/PulsarCommandSender.java | 8 ++++ .../service/PulsarCommandSenderImpl.java | 45 +++++++++++++++++++ .../pulsar/broker/service/ServerCnx.java | 24 +++------- .../intercept/CounterBrokerInterceptor.java | 30 +++++++++++++ .../broker/transaction/TransactionTest.java | 10 ++++- .../pulsar/common/protocol/Commands.java | 17 +++---- 9 files changed, 117 insertions(+), 37 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptor.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptor.java index ee2cc6185e0b9..1c4c09f8f2566 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptor.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptor.java @@ -129,7 +129,7 @@ default void messageAcked(ServerCnx cnx, Consumer consumer, * @param tcId Transaction Coordinator Id * @param txnID Transaction ID */ - default void beginTxn(long tcId, String txnID) { + default void txnOpened(long tcId, String txnID) { } /** @@ -138,7 +138,7 @@ default void beginTxn(long tcId, String txnID) { * @param txnID Transaction ID * @param txnAction Transaction Action */ - default void endTxn(String txnID, long txnAction) { + default void txnEnded(String txnID, long txnAction) { } /** * Called by the broker while new command incoming. diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java index 82039aa4a90e5..f04446fa9a06d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java @@ -106,13 +106,13 @@ public void messageAcked(ServerCnx cnx, Consumer consumer, } @Override - public void beginTxn(long tcId, String txnID) { - this.interceptor.beginTxn(tcId, txnID); + public void txnOpened(long tcId, String txnID) { + this.interceptor.txnOpened(tcId, txnID); } @Override - public void endTxn(String txnID, long txnAction) { - this.interceptor.endTxn(txnID, txnAction); + public void txnEnded(String txnID, long txnAction) { + this.interceptor.txnEnded(txnID, txnAction); } @Override diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptors.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptors.java index 75211276f8e89..225066b94342f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptors.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptors.java @@ -167,22 +167,22 @@ public void messageAcked(ServerCnx cnx, Consumer consumer, } @Override - public void beginTxn(long tcId, String txnID) { + public void txnOpened(long tcId, String txnID) { if (interceptors == null || interceptors.isEmpty()) { return; } for (BrokerInterceptorWithClassLoader value : interceptors.values()) { - value.beginTxn(tcId, txnID); + value.txnOpened(tcId, txnID); } } @Override - public void endTxn(String txnID, long txnAction) { + public void txnEnded(String txnID, long txnAction) { if (interceptors == null || interceptors.isEmpty()) { return; } for (BrokerInterceptorWithClassLoader value : interceptors.values()) { - value.endTxn(txnID, txnAction); + value.txnEnded(txnID, txnAction); } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarCommandSender.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarCommandSender.java index 0ecda2d012426..2032d96bf25dc 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarCommandSender.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarCommandSender.java @@ -23,6 +23,7 @@ import java.util.List; import java.util.Optional; import org.apache.bookkeeper.mledger.Entry; +import org.apache.pulsar.client.api.transaction.TxnID; import org.apache.pulsar.common.api.proto.CommandLookupTopicResponse; import org.apache.pulsar.common.api.proto.ServerError; import org.apache.pulsar.common.protocol.schema.SchemaVersion; @@ -83,4 +84,11 @@ Future sendMessagesToConsumer(long consumerId, String topicName, Subscript void sendTcClientConnectResponse(long requestId); + void sendNewTxnResponse(long requestId, TxnID txnID, long tcID); + + void sendNewTxnErrorResponse(long requestId, long txnID, ServerError error, String message); + + void sendEndTxnResponse(long requestId, TxnID txnID, int txnAction); + + void sendEndTxnErrorResponse(long requestId, TxnID txnID, ServerError error, String message); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarCommandSenderImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarCommandSenderImpl.java index 7c1a920482145..27d2e58c9dff7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarCommandSenderImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarCommandSenderImpl.java @@ -28,10 +28,12 @@ import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.pulsar.broker.intercept.BrokerInterceptor; +import org.apache.pulsar.client.api.transaction.TxnID; import org.apache.pulsar.common.api.proto.BaseCommand; import org.apache.pulsar.common.api.proto.CommandLookupTopicResponse; import org.apache.pulsar.common.api.proto.ProtocolVersion; import org.apache.pulsar.common.api.proto.ServerError; +import org.apache.pulsar.common.api.proto.TxnAction; import org.apache.pulsar.common.protocol.Commands; import org.apache.pulsar.common.protocol.schema.SchemaVersion; import org.apache.pulsar.common.schema.SchemaInfo; @@ -299,6 +301,49 @@ public void sendTcClientConnectResponse(long requestId) { sendTcClientConnectResponse(requestId, null, null); } + @Override + public void sendNewTxnResponse(long requestId, TxnID txnID, long tcID) { + BaseCommand command = Commands.newTxnResponse(requestId, txnID.getLeastSigBits(), + txnID.getMostSigBits()); + safeIntercept(command, cnx); + ByteBuf outBuf = Commands.serializeWithSize(command); + cnx.ctx().writeAndFlush(outBuf); + if (this.interceptor != null) { + this.interceptor.txnOpened(tcID, txnID.toString()); + } + } + + @Override + public void sendNewTxnErrorResponse(long requestId, long txnID, ServerError error, String message) { + BaseCommand command = Commands.newTxnResponse(requestId, txnID, error, message); + safeIntercept(command, cnx); + ByteBuf outBuf = Commands.serializeWithSize(command); + cnx.ctx().writeAndFlush(outBuf); + } + + @Override + public void sendEndTxnResponse(long requestId, TxnID txnID, int txnAction) { + BaseCommand command = Commands.newEndTxnResponse(requestId, txnID.getLeastSigBits(), + txnID.getMostSigBits()); + safeIntercept(command, cnx); + ByteBuf outBuf = Commands.serializeWithSize(command); + cnx.ctx().writeAndFlush(outBuf); + if (this.interceptor != null) { + this.interceptor.txnEnded(txnID.toString(), txnAction); + } + } + + @Override + public void sendEndTxnErrorResponse(long requestId, TxnID txnID, ServerError error, String message) { + BaseCommand command = Commands.newEndTxnResponse(requestId, txnID.getMostSigBits(), error, message); + safeIntercept(command, cnx); + ByteBuf outBuf = Commands.serializeWithSize(command); + cnx.ctx().writeAndFlush(outBuf); + if (this.interceptor != null) { + this.interceptor.txnEnded(txnID.toString(), TxnAction.ABORT_VALUE); + } + } + private void safeIntercept(BaseCommand command, ServerCnx cnx) { try { this.interceptor.onPulsarCommand(command, cnx); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index 49454acf185d2..7ececa6d62f5f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -2175,16 +2175,12 @@ protected void handleNewTxn(CommandNewTxn command) { if (log.isDebugEnabled()) { log.debug("Send response {} for new txn request {}", tcId.getId(), requestId); } - ctx.writeAndFlush(Commands.newTxnResponse(requestId, txnID.getLeastSigBits(), - txnID.getMostSigBits())); - if (getBrokerService().getInterceptor() != null) { - getBrokerService().getInterceptor().beginTxn(command.getTcId(), txnID.toString()); - } + commandSender.sendNewTxnResponse(requestId, txnID, command.getTcId()); } else { ex = handleTxnException(ex, BaseCommand.Type.NEW_TXN.name(), requestId); - ctx.writeAndFlush(Commands.newTxnResponse(requestId, tcId.getId(), - BrokerServiceException.getClientErrorCode(ex), ex.getMessage())); + commandSender.sendNewTxnErrorResponse(requestId, tcId.getId(), + BrokerServiceException.getClientErrorCode(ex), ex.getMessage()); transactionMetadataStoreService.handleOpFail(ex, tcId); } })); @@ -2245,20 +2241,14 @@ protected void handleEndTxn(CommandEndTxn command) { .endTransaction(txnID, txnAction, false) .whenComplete((v, ex) -> { if (ex == null) { - ctx.writeAndFlush(Commands.newEndTxnResponse(requestId, - txnID.getLeastSigBits(), txnID.getMostSigBits())); - if (getBrokerService().getInterceptor() != null) { - getBrokerService().getInterceptor().endTxn(txnID.toString(), txnAction); - } + commandSender.sendEndTxnResponse(requestId, + txnID, txnAction); } else { ex = handleTxnException(ex, BaseCommand.Type.END_TXN.name(), requestId); - ctx.writeAndFlush(Commands.newEndTxnResponse(requestId, txnID.getMostSigBits(), - BrokerServiceException.getClientErrorCode(ex), ex.getMessage())); + commandSender.sendEndTxnErrorResponse(requestId, txnID, + BrokerServiceException.getClientErrorCode(ex), ex.getMessage()); transactionMetadataStoreService.handleOpFail(ex, tcId); - if (getBrokerService().getInterceptor() != null) { - getBrokerService().getInterceptor().endTxn(txnID.toString(), TxnAction.ABORT_VALUE); - } } }); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/CounterBrokerInterceptor.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/CounterBrokerInterceptor.java index 991f211024570..4ae4cf00acb01 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/CounterBrokerInterceptor.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/CounterBrokerInterceptor.java @@ -41,6 +41,7 @@ import org.apache.pulsar.common.api.proto.BaseCommand; import org.apache.pulsar.common.api.proto.MessageMetadata; import org.apache.pulsar.common.api.proto.CommandAck; +import org.apache.pulsar.common.api.proto.TxnAction; import org.eclipse.jetty.server.Response; @@ -56,6 +57,9 @@ public class CounterBrokerInterceptor implements BrokerInterceptor { int messageDispatchCount = 0; int messageAckCount = 0; int handleAckCount = 0; + int txnCount = 0; + int committedTxnCount = 0; + int abortedTxnCount = 0; private List responseList = new ArrayList<>(); @@ -177,6 +181,20 @@ public void onFilter(ServletRequest request, ServletResponse response, FilterCha chain.doFilter(request, response); } + @Override + public void txnOpened(long tcId, String txnID) { + txnCount ++; + } + + @Override + public void txnEnded(String txnID, long txnAction) { + if(txnAction == TxnAction.COMMIT_VALUE) { + committedTxnCount ++; + } else { + abortedTxnCount ++; + } + } + @Override public void initialize(PulsarService pulsarService) throws Exception { @@ -230,4 +248,16 @@ public void clearResponseList() { public List getResponseList() { return responseList; } + + public int getTxnCount() { + return txnCount; + } + + public int getCommittedTxnCount() { + return committedTxnCount; + } + + public int getAbortedTxnCount() { + return abortedTxnCount; + } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java index 07a7c951e9ca2..b9a72b3e40a44 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java @@ -60,6 +60,8 @@ import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.pulsar.broker.PulsarService; +import org.apache.pulsar.broker.intercept.BrokerInterceptor; +import org.apache.pulsar.broker.intercept.CounterBrokerInterceptor; import org.apache.pulsar.broker.service.Topic; import org.apache.pulsar.broker.service.persistent.PersistentSubscription; import org.apache.pulsar.broker.service.persistent.PersistentTopic; @@ -699,8 +701,12 @@ public void testEndTxnWhenCommittingOrAborting() throws Exception { field.set(commitTxn, TransactionImpl.State.COMMITTING); field.set(abortTxn, TransactionImpl.State.ABORTING); - abortTxn.abort(); - commitTxn.commit(); + BrokerInterceptor listener = getPulsarServiceList().get(0).getBrokerInterceptor(); + assertEquals(((CounterBrokerInterceptor)listener).getTxnCount(),2); + abortTxn.abort().get(); + assertEquals(((CounterBrokerInterceptor)listener).getAbortedTxnCount(),1); + commitTxn.commit().get(); + assertEquals(((CounterBrokerInterceptor)listener).getCommittedTxnCount(),1); } @Test diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java index 95b21fdc9e493..85843c800910c 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java @@ -1264,16 +1264,16 @@ public static ByteBuf newTxn(long tcId, long requestId, long ttlSeconds) { return serializeWithSize(cmd); } - public static ByteBuf newTxnResponse(long requestId, long txnIdLeastBits, long txnIdMostBits) { + public static BaseCommand newTxnResponse(long requestId, long txnIdLeastBits, long txnIdMostBits) { BaseCommand cmd = localCmd(Type.NEW_TXN_RESPONSE); cmd.setNewTxnResponse() .setRequestId(requestId) .setTxnidMostBits(txnIdMostBits) .setTxnidLeastBits(txnIdLeastBits); - return serializeWithSize(cmd); + return cmd; } - public static ByteBuf newTxnResponse(long requestId, long txnIdMostBits, ServerError error, String errorMsg) { + public static BaseCommand newTxnResponse(long requestId, long txnIdMostBits, ServerError error, String errorMsg) { BaseCommand cmd = localCmd(Type.NEW_TXN_RESPONSE); CommandNewTxnResponse response = cmd.setNewTxnResponse() .setRequestId(requestId) @@ -1282,7 +1282,7 @@ public static ByteBuf newTxnResponse(long requestId, long txnIdMostBits, ServerE if (errorMsg != null) { response.setMessage(errorMsg); } - return serializeWithSize(cmd); + return cmd; } public static ByteBuf newAddPartitionToTxn(long requestId, long txnIdLeastBits, long txnIdMostBits, @@ -1363,16 +1363,17 @@ public static BaseCommand newEndTxn(long requestId, long txnIdLeastBits, long tx return cmd; } - public static ByteBuf newEndTxnResponse(long requestId, long txnIdLeastBits, long txnIdMostBits) { + public static BaseCommand newEndTxnResponse(long requestId, long txnIdLeastBits, long txnIdMostBits) { BaseCommand cmd = localCmd(Type.END_TXN_RESPONSE); cmd.setEndTxnResponse() .setRequestId(requestId) .setTxnidLeastBits(txnIdLeastBits) .setTxnidMostBits(txnIdMostBits); - return serializeWithSize(cmd); + return cmd; } - public static ByteBuf newEndTxnResponse(long requestId, long txnIdMostBits, ServerError error, String errorMsg) { + public static BaseCommand newEndTxnResponse(long requestId, long txnIdMostBits, + ServerError error, String errorMsg) { BaseCommand cmd = localCmd(Type.END_TXN_RESPONSE); CommandEndTxnResponse response = cmd.setEndTxnResponse() .setRequestId(requestId) @@ -1381,7 +1382,7 @@ public static ByteBuf newEndTxnResponse(long requestId, long txnIdMostBits, Serv if (errorMsg != null) { response.setMessage(errorMsg); } - return serializeWithSize(cmd); + return cmd; } public static ByteBuf newEndTxnOnPartition(long requestId, long txnIdLeastBits, long txnIdMostBits, String topic,