From 0a753d75b9ef2c2263d4c6467a7e1391ae23a53c Mon Sep 17 00:00:00 2001 From: penghui Date: Thu, 11 Aug 2022 00:44:29 +0800 Subject: [PATCH 01/10] [fix][broker] Fix schema does not replicate successfully ### Motivation #11441 supports replicate schema to remote clusters. But there is a mistake that the returned schema state is incorrect. https://github.com/apache/pulsar/blob/e826d849ceef9d6aef28569ad57950bba90dfff1/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessageImpl.java#L765-L770 Because the replicator used MessageImpl will not have the schema. And this will cause the producer to skip the schema upload. https://github.com/apache/pulsar/blob/e826d849ceef9d6aef28569ad57950bba90dfff1/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java#L2147-L2149 We should remove https://github.com/apache/pulsar/blob/e826d849ceef9d6aef28569ad57950bba90dfff1/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessageImpl.java#L766-L768 To return the correct schema state. And then we should also provide the correct schema hash. If the message is used by the replicator, the schema hash should be based on the replicator schema. Otherwise, it should use based on the schema of the message. ### Modification - Fixed the incorrect returned schema state - Provide the method for getting schema hash for MessageImpl ### Verification Update the test only to create producer to one cluster. Because if create a producer for other clusters, the producer will upload the schema. This is the reason that why the test can get pass before. --- .../pulsar/broker/service/ReplicatorTest.java | 29 +++++++++++-------- .../pulsar/client/impl/MessageImpl.java | 14 +++++++-- .../pulsar/client/impl/ProducerImpl.java | 9 ++---- .../common/protocol/schema/SchemaHash.java | 4 +++ 4 files changed, 35 insertions(+), 21 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java index c8c257992cf02..fd75bdfdba640 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java @@ -70,6 +70,7 @@ import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; +import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.MessageRoutingMode; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.PulsarClient; @@ -392,19 +393,24 @@ public void testReplicationWithSchema() throws Exception { final String subName = "my-sub"; @Cleanup - Producer producer1 = client1.newProducer(Schema.AVRO(Schemas.PersonOne.class)) - .topic(topic.toString()) - .create(); - @Cleanup - Producer producer2 = client2.newProducer(Schema.AVRO(Schemas.PersonOne.class)) - .topic(topic.toString()) - .create(); - @Cleanup - Producer producer3 = client3.newProducer(Schema.AVRO(Schemas.PersonOne.class)) + Producer producer = client1.newProducer(Schema.AVRO(Schemas.PersonOne.class)) .topic(topic.toString()) .create(); - List> producers = Lists.newArrayList(producer1, producer2, producer3); + admin1.topics().createSubscription(topic.toString(), subName, MessageId.earliest); + admin2.topics().createSubscription(topic.toString(), subName, MessageId.earliest); + admin3.topics().createSubscription(topic.toString(), subName, MessageId.earliest); + + + for (int i = 0; i < 10; i++) { + producer.send(new Schemas.PersonOne(i)); + } + + Awaitility.await().untilAsserted(() -> { + assertTrue(admin1.topics().getInternalStats(topic.toString()).schemaLedgers.size() > 0); + assertTrue(admin2.topics().getInternalStats(topic.toString()).schemaLedgers.size() > 0); + assertTrue(admin3.topics().getInternalStats(topic.toString()).schemaLedgers.size() > 0); + }); @Cleanup Consumer consumer1 = client1.newConsumer(Schema.AVRO(Schemas.PersonOne.class)) @@ -424,8 +430,7 @@ public void testReplicationWithSchema() throws Exception { .subscriptionName(subName) .subscribe(); - for (int i = 0; i < 3; i++) { - producers.get(i).send(new Schemas.PersonOne(i)); + for (int i = 0; i < 10; i++) { Message msg1 = consumer1.receive(); Message msg2 = consumer2.receive(); Message msg3 = consumer3.receive(); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessageImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessageImpl.java index 4970f95b7fc69..cee30148f4448 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessageImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessageImpl.java @@ -52,6 +52,7 @@ import org.apache.pulsar.common.api.proto.SingleMessageMetadata; import org.apache.pulsar.common.protocol.Commands; import org.apache.pulsar.common.protocol.schema.BytesSchemaVersion; +import org.apache.pulsar.common.protocol.schema.SchemaHash; import org.apache.pulsar.common.schema.KeyValueEncodingType; import org.apache.pulsar.common.schema.SchemaInfo; import org.apache.pulsar.common.schema.SchemaType; @@ -64,6 +65,8 @@ public class MessageImpl implements Message { private ByteBuf payload; private Schema schema; + + private SchemaHash schemaHash; private SchemaInfo schemaInfoForReplicator; private SchemaState schemaState = SchemaState.None; private Optional encryptionCtx = Optional.empty(); @@ -91,6 +94,9 @@ public static MessageImpl create(MessageMetadata msgMetadata, ByteBuffer msg.payload = Unpooled.wrappedBuffer(payload); msg.properties = null; msg.schema = schema; + if (msg.schema != null) { + msg.schemaHash = SchemaHash.of(schema); + } msg.uncompressedSize = payload.remaining(); return msg; } @@ -431,9 +437,14 @@ public SchemaInfo getSchemaInfo() { return schema.getSchemaInfo(); } + public SchemaHash getSchemaHash() { + return schemaHash; + } + public void setSchemaInfoForReplicator(SchemaInfo schemaInfo) { if (msgMetadata.hasReplicatedFrom()) { this.schemaInfoForReplicator = schemaInfo; + this.schemaHash = SchemaHash.of(schemaInfo); } else { throw new IllegalArgumentException("Only allowed to set schemaInfoForReplicator for a replicated message."); } @@ -763,9 +774,6 @@ int getUncompressedSize() { } SchemaState getSchemaState() { - if (getSchemaInfo() == null) { - return SchemaState.Ready; - } return schemaState; } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index b5756e61b2c7a..6a11617968c65 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -717,8 +717,7 @@ private boolean populateMessageSchema(MessageImpl msg, SendCallback callback) { completeCallbackAndReleaseSemaphore(msg.getUncompressedSize(), callback, e); return false; } - SchemaHash schemaHash = SchemaHash.of(msg.getSchemaInternal()); - byte[] schemaVersion = schemaCache.get(schemaHash); + byte[] schemaVersion = schemaCache.get(msg.getSchemaHash()); if (schemaVersion != null) { msgMetadataBuilder.setSchemaVersion(schemaVersion); msg.setSchemaState(MessageImpl.SchemaState.Ready); @@ -727,8 +726,7 @@ private boolean populateMessageSchema(MessageImpl msg, SendCallback callback) { } private boolean rePopulateMessageSchema(MessageImpl msg) { - SchemaHash schemaHash = SchemaHash.of(msg.getSchemaInternal()); - byte[] schemaVersion = schemaCache.get(schemaHash); + byte[] schemaVersion = schemaCache.get(msg.getSchemaHash()); if (schemaVersion == null) { return false; } @@ -759,8 +757,7 @@ private void tryRegisterSchema(ClientCnx cnx, MessageImpl msg, SendCallback call // case, we should not cache the schema version so that the schema version of the message metadata will // be null, instead of an empty array. if (v.length != 0) { - SchemaHash schemaHash = SchemaHash.of(msg.getSchemaInternal()); - schemaCache.putIfAbsent(schemaHash, v); + schemaCache.putIfAbsent(msg.getSchemaHash(), v); msg.getMessageBuilder().setSchemaVersion(v); } msg.setSchemaState(MessageImpl.SchemaState.Ready); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/schema/SchemaHash.java b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/schema/SchemaHash.java index 40220e6047a3b..228575bd170e0 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/schema/SchemaHash.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/schema/SchemaHash.java @@ -54,6 +54,10 @@ public static SchemaHash of(SchemaData schemaData) { return of(schemaData.getData(), schemaData.getType()); } + public static SchemaHash of(SchemaInfo schemaInfo) { + return of(schemaInfo.getSchema(), schemaInfo.getType()); + } + private static SchemaHash of(byte[] schemaBytes, SchemaType schemaType) { return new SchemaHash(hashFunction.hashBytes(schemaBytes), schemaType); } From 691fd253e0d23093477170c5cdf3c4c5206751f5 Mon Sep 17 00:00:00 2001 From: penghui Date: Thu, 11 Aug 2022 11:44:49 +0800 Subject: [PATCH 02/10] Fix NPE --- .../org/apache/pulsar/common/protocol/schema/SchemaHash.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/schema/SchemaHash.java b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/schema/SchemaHash.java index 228575bd170e0..73bfde6455c63 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/schema/SchemaHash.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/schema/SchemaHash.java @@ -55,7 +55,8 @@ public static SchemaHash of(SchemaData schemaData) { } public static SchemaHash of(SchemaInfo schemaInfo) { - return of(schemaInfo.getSchema(), schemaInfo.getType()); + return of(schemaInfo == null ? new byte[0] : schemaInfo.getSchema(), + schemaInfo == null ? null : schemaInfo.getType()); } private static SchemaHash of(byte[] schemaBytes, SchemaType schemaType) { From a75661349f1d646e2ff5d9cc444c2e3feb2f75cb Mon Sep 17 00:00:00 2001 From: penghui Date: Thu, 11 Aug 2022 16:14:36 +0800 Subject: [PATCH 03/10] Fix NPE --- .../java/org/apache/pulsar/client/impl/MessageImpl.java | 6 ++---- .../apache/pulsar/common/protocol/schema/SchemaHash.java | 2 +- 2 files changed, 3 insertions(+), 5 deletions(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessageImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessageImpl.java index cee30148f4448..6d1d586b49b7f 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessageImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessageImpl.java @@ -94,9 +94,7 @@ public static MessageImpl create(MessageMetadata msgMetadata, ByteBuffer msg.payload = Unpooled.wrappedBuffer(payload); msg.properties = null; msg.schema = schema; - if (msg.schema != null) { - msg.schemaHash = SchemaHash.of(schema); - } + msg.schemaHash = SchemaHash.of(schema); msg.uncompressedSize = payload.remaining(); return msg; } @@ -438,7 +436,7 @@ public SchemaInfo getSchemaInfo() { } public SchemaHash getSchemaHash() { - return schemaHash; + return schemaHash == null ? SchemaHash.of(new byte[0], null) : schemaHash; } public void setSchemaInfoForReplicator(SchemaInfo schemaInfo) { diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/schema/SchemaHash.java b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/schema/SchemaHash.java index 73bfde6455c63..8bbc18fbb703c 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/schema/SchemaHash.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/schema/SchemaHash.java @@ -59,7 +59,7 @@ public static SchemaHash of(SchemaInfo schemaInfo) { schemaInfo == null ? null : schemaInfo.getType()); } - private static SchemaHash of(byte[] schemaBytes, SchemaType schemaType) { + public static SchemaHash of(byte[] schemaBytes, SchemaType schemaType) { return new SchemaHash(hashFunction.hashBytes(schemaBytes), schemaType); } From a9958832125905a4d07c4581d88adfcbf78e7378 Mon Sep 17 00:00:00 2001 From: penghui Date: Wed, 17 Aug 2022 00:20:52 +0800 Subject: [PATCH 04/10] Fix Test --- .../pulsar/broker/service/ReplicatorTest.java | 25 ++++++++++--------- 1 file changed, 13 insertions(+), 12 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java index fd75bdfdba640..96bc8258a7d14 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java @@ -62,6 +62,7 @@ import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl.State; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; +import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.pulsar.broker.BrokerTestUtil; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.service.BrokerServiceException.NamingException; @@ -79,6 +80,7 @@ import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.client.api.TypedMessageBuilder; +import org.apache.pulsar.client.impl.MessageIdImpl; import org.apache.pulsar.client.impl.ProducerImpl; import org.apache.pulsar.client.impl.PulsarClientImpl; import org.apache.pulsar.client.impl.conf.ProducerConfigurationData; @@ -1400,6 +1402,8 @@ public void testReplicatorWithFailedAck() throws Exception { PersistentTopic topic = (PersistentTopic) pulsar1.getBrokerService().getTopic(dest.toString(), false) .getNow(null).get(); + MessageIdImpl lastMessageId = (MessageIdImpl) topic.getLastMessageId().get(); + Position lastPosition = PositionImpl.get(lastMessageId.getLedgerId(), lastMessageId.getEntryId()); ConcurrentOpenHashMap replicators = topic.getReplicators(); PersistentReplicator replicator = (PersistentReplicator) replicators.get("r2"); @@ -1409,6 +1413,9 @@ public void testReplicatorWithFailedAck() throws Exception { assertEquals(replicator.getState(), org.apache.pulsar.broker.service.AbstractReplicator.State.Started); ManagedCursorImpl cursor = (ManagedCursorImpl) replicator.getCursor(); + + // Make sure all the data has replicated to the remote cluster before close the cursor. + Awaitility.await().untilAsserted(() -> assertEquals(cursor.getMarkDeletedPosition(), lastPosition)); cursor.setState(State.Closed); Field field = ManagedCursorImpl.class.getDeclaredField("state"); @@ -1417,22 +1424,16 @@ public void testReplicatorWithFailedAck() throws Exception { producer1.produce(10); - Position deletedPos = cursor.getMarkDeletedPosition(); - Position readPos = cursor.getReadPosition(); - - Awaitility.await().timeout(30, TimeUnit.SECONDS).until( - () -> cursor.getMarkDeletedPosition().getEntryId() != (cursor.getReadPosition().getEntryId() - 1)); - - assertNotEquals((readPos.getEntryId() - 1), deletedPos.getEntryId()); + // The cursor is closed, so the mark delete position will not move forward. + assertEquals(cursor.getMarkDeletedPosition(), lastPosition); field.set(cursor, State.Open); Awaitility.await().timeout(30, TimeUnit.SECONDS).until( - () -> cursor.getMarkDeletedPosition().getEntryId() == (cursor.getReadPosition().getEntryId() - 1)); - - deletedPos = cursor.getMarkDeletedPosition(); - readPos = cursor.getReadPosition(); - assertEquals((readPos.getEntryId() - 1), deletedPos.getEntryId()); + () -> { + log.info("++++++++++++ {}, {}", cursor.getMarkDeletedPosition(), cursor.getReadPosition()); + return cursor.getMarkDeletedPosition().getEntryId() == (cursor.getReadPosition().getEntryId() - 1); + }); } private static final Logger log = LoggerFactory.getLogger(ReplicatorTest.class); From ee1563f8e28db1e2336a28654a405e2f130ab892 Mon Sep 17 00:00:00 2001 From: penghui Date: Wed, 17 Aug 2022 16:48:22 +0800 Subject: [PATCH 05/10] Fix Test --- .../org/apache/pulsar/client/impl/ConnectionTimeoutTest.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConnectionTimeoutTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConnectionTimeoutTest.java index af79a4f645f2d..e76ae0d897cce 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConnectionTimeoutTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConnectionTimeoutTest.java @@ -18,7 +18,6 @@ */ package org.apache.pulsar.client.impl; -import io.netty.channel.ConnectTimeoutException; import java.io.IOException; import java.net.InetAddress; import java.net.ServerSocket; @@ -73,7 +72,9 @@ public void testLowTimeout() throws Exception { } catch (TimeoutException e) { Assert.fail("Connection timeout didn't apply."); } catch (Exception e) { - Assert.assertEquals(e.getCause().getCause().getCause().getClass(), ConnectTimeoutException.class); + String causeClassName = e.getCause().getCause().getCause().getClass().getName(); + Assert.assertTrue(causeClassName.equals("io.netty.channel.ConnectTimeoutException") + || causeClassName.equals("io.netty.channel.StacklessClosedChannelException")); } } finally { threads.stream().forEach(Thread::interrupt); From 221e74b14c044e45b030d3f77dfb179939241cde Mon Sep 17 00:00:00 2001 From: penghui Date: Wed, 17 Aug 2022 17:06:39 +0800 Subject: [PATCH 06/10] Revert "Fix Test" This reverts commit ee1563f8e28db1e2336a28654a405e2f130ab892. --- .../org/apache/pulsar/client/impl/ConnectionTimeoutTest.java | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConnectionTimeoutTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConnectionTimeoutTest.java index e76ae0d897cce..af79a4f645f2d 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConnectionTimeoutTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConnectionTimeoutTest.java @@ -18,6 +18,7 @@ */ package org.apache.pulsar.client.impl; +import io.netty.channel.ConnectTimeoutException; import java.io.IOException; import java.net.InetAddress; import java.net.ServerSocket; @@ -72,9 +73,7 @@ public void testLowTimeout() throws Exception { } catch (TimeoutException e) { Assert.fail("Connection timeout didn't apply."); } catch (Exception e) { - String causeClassName = e.getCause().getCause().getCause().getClass().getName(); - Assert.assertTrue(causeClassName.equals("io.netty.channel.ConnectTimeoutException") - || causeClassName.equals("io.netty.channel.StacklessClosedChannelException")); + Assert.assertEquals(e.getCause().getCause().getCause().getClass(), ConnectTimeoutException.class); } } finally { threads.stream().forEach(Thread::interrupt); From 312b4d5e5e01a2e5bc8596d69db811ba744ecde7 Mon Sep 17 00:00:00 2001 From: penghui Date: Wed, 17 Aug 2022 17:07:01 +0800 Subject: [PATCH 07/10] Disable reuseFork --- build/run_unit_group.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/build/run_unit_group.sh b/build/run_unit_group.sh index 4804c236efc6e..bc8382dc0eb16 100755 --- a/build/run_unit_group.sh +++ b/build/run_unit_group.sh @@ -134,7 +134,7 @@ function test_group_other() { -Dexclude='**/ManagedLedgerTest.java, **/OffloadersCacheTest.java **/PrimitiveSchemaTest.java, - BlobStoreManagedLedgerOffloaderTest.java' + BlobStoreManagedLedgerOffloaderTest.java' -DreuseForks=false mvn_test -pl managed-ledger -Dinclude='**/ManagedLedgerTest.java, **/OffloadersCacheTest.java' From 98eb1d4e35c3cb0c5d6e94a4dd610368804c110b Mon Sep 17 00:00:00 2001 From: Jiwe Guo Date: Tue, 16 Aug 2022 15:20:50 +0800 Subject: [PATCH 08/10] fix test. --- .../java/org/apache/pulsar/broker/service/ReplicatorTest.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java index 96bc8258a7d14..be232760b0b80 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java @@ -1407,7 +1407,7 @@ public void testReplicatorWithFailedAck() throws Exception { ConcurrentOpenHashMap replicators = topic.getReplicators(); PersistentReplicator replicator = (PersistentReplicator) replicators.get("r2"); - Awaitility.await().timeout(50, TimeUnit.SECONDS) + Awaitility.await().pollInterval(1, TimeUnit.SECONDS).timeout(30, TimeUnit.SECONDS) .untilAsserted(() -> assertEquals(org.apache.pulsar.broker.service.AbstractReplicator.State.Started, replicator.getState())); @@ -1416,6 +1416,7 @@ public void testReplicatorWithFailedAck() throws Exception { // Make sure all the data has replicated to the remote cluster before close the cursor. Awaitility.await().untilAsserted(() -> assertEquals(cursor.getMarkDeletedPosition(), lastPosition)); + cursor.setState(State.Closed); Field field = ManagedCursorImpl.class.getDeclaredField("state"); From 0e3d62bcfde3407dc47850549c3dff50c84864d7 Mon Sep 17 00:00:00 2001 From: Jiwe Guo Date: Thu, 18 Aug 2022 09:49:45 +0800 Subject: [PATCH 09/10] fix param --- build/run_unit_group.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/build/run_unit_group.sh b/build/run_unit_group.sh index bc8382dc0eb16..87836261609b2 100755 --- a/build/run_unit_group.sh +++ b/build/run_unit_group.sh @@ -134,7 +134,7 @@ function test_group_other() { -Dexclude='**/ManagedLedgerTest.java, **/OffloadersCacheTest.java **/PrimitiveSchemaTest.java, - BlobStoreManagedLedgerOffloaderTest.java' -DreuseForks=false + BlobStoreManagedLedgerOffloaderTest.java' -DtestReuseFork=false mvn_test -pl managed-ledger -Dinclude='**/ManagedLedgerTest.java, **/OffloadersCacheTest.java' From cdd89cf321bb8dd5561344908f2b74b42bb82c59 Mon Sep 17 00:00:00 2001 From: Jiwe Guo Date: Thu, 18 Aug 2022 09:51:15 +0800 Subject: [PATCH 10/10] fix test. --- .../apache/pulsar/client/impl/MultiTopicsConsumerImplTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImplTest.java index 51fba75cd216b..15288cecb1ca2 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImplTest.java @@ -126,7 +126,7 @@ public void testGetStats() throws Exception { // // Code under tests is using CompletableFutures. Theses may hang indefinitely if code is broken. // That's why a test timeout is defined. - @Test(timeOut = 5000) + @Test(timeOut = 10000) public void testParallelSubscribeAsync() throws Exception { String topicName = "parallel-subscribe-async-topic"; MultiTopicsConsumerImpl impl = createMultiTopicsConsumer();