diff --git a/pulsar-broker-common/src/test/java/org/apache/pulsar/common/configuration/PulsarConfigurationLoaderTest.java b/pulsar-broker-common/src/test/java/org/apache/pulsar/common/configuration/PulsarConfigurationLoaderTest.java index 7e37a08a3f70f..9714cbc7c7e8b 100644 --- a/pulsar-broker-common/src/test/java/org/apache/pulsar/common/configuration/PulsarConfigurationLoaderTest.java +++ b/pulsar-broker-common/src/test/java/org/apache/pulsar/common/configuration/PulsarConfigurationLoaderTest.java @@ -115,7 +115,7 @@ public void testPulsarConfiguraitonLoadingStream() throws Exception { final ServiceConfiguration serviceConfig = PulsarConfigurationLoader.create(stream, ServiceConfiguration.class); assertNotNull(serviceConfig); assertEquals(serviceConfig.getMetadataStoreUrl(), zkServer); - assertEquals(serviceConfig.isBrokerDeleteInactiveTopicsEnabled(), true); + assertTrue(serviceConfig.isBrokerDeleteInactiveTopicsEnabled()); assertEquals(serviceConfig.getBacklogQuotaDefaultLimitGB(), 18); assertEquals(serviceConfig.getClusterName(), "usc"); assertEquals(serviceConfig.getBrokerClientAuthenticationParameters(), "role:my-role"); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarSinkE2ETest.java b/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarSinkE2ETest.java index e82dc03b83fc9..8cb581669f8c1 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarSinkE2ETest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarSinkE2ETest.java @@ -118,7 +118,7 @@ public void testReadCompactedSink() throws Exception { consumerProperties.put("readCompacted","true"); sinkConfig.setInputSpecs(Collections.singletonMap(sourceTopic, ConsumerConfig.builder().consumerProperties(consumerProperties).build())); String jarFilePathUrl = getPulsarIODataGeneratorNar().toURI().toString(); - admin.sink().createSinkWithUrl(sinkConfig, jarFilePathUrl); + admin.sinks().createSinkWithUrl(sinkConfig, jarFilePathUrl); // 5 Sink should only read compacted value,so we will only receive compacted messages Awaitility.await().ignoreExceptions().untilAsserted(() -> { diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/MessageRoutingMode.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/MessageRoutingMode.java index f37fa93542fec..4f66647645163 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/MessageRoutingMode.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/MessageRoutingMode.java @@ -24,7 +24,7 @@ /** * Default routing mode for messages to partition. * - *

This logic is applied when the application is not setting a key {@link MessageBuilder#setKey(String)} + *

This logic is applied when the application is not setting a key {@link TypedMessageBuilder#key(String)} * on a particular message. */ @InterfaceAudience.Public diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/stats/JvmMetrics.java b/pulsar-common/src/main/java/org/apache/pulsar/common/stats/JvmMetrics.java index 83b14060a99a2..a35760e53c34b 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/stats/JvmMetrics.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/stats/JvmMetrics.java @@ -121,7 +121,7 @@ public List generate() { long totalAllocated = 0; long totalUsed = 0; - for (PoolArenaMetric arena : PooledByteBufAllocator.DEFAULT.directArenas()) { + for (PoolArenaMetric arena : PooledByteBufAllocator.DEFAULT.metric().directArenas()) { this.gcLogger.logMetrics(m); for (PoolChunkListMetric list : arena.chunkLists()) { for (PoolChunkMetric chunk : list) { diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SinksImpl.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SinksImpl.java index ab69ab9182cf9..6d867b37d9307 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SinksImpl.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SinksImpl.java @@ -524,7 +524,7 @@ public SinkStatus getStatus(final String tenant, sinkInstanceStatusData = getComponentInstanceStatus(tenant, namespace, name, assignment.getInstance().getInstanceId(), null); } else { - sinkInstanceStatusData = worker().getFunctionAdmin().sink().getSinkStatus( + sinkInstanceStatusData = worker().getFunctionAdmin().sinks().getSinkStatus( assignment.getInstance().getFunctionMetaData().getFunctionDetails().getTenant(), assignment.getInstance().getFunctionMetaData().getFunctionDetails().getNamespace(), assignment.getInstance().getFunctionMetaData().getFunctionDetails().getName(), diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SourcesImpl.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SourcesImpl.java index df2dca813e770..e3946ac3642ae 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SourcesImpl.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SourcesImpl.java @@ -533,7 +533,7 @@ public SourceStatus getStatus(final String tenant, if (isOwner) { sourceInstanceStatusData = getComponentInstanceStatus(tenant, namespace, name, assignment.getInstance().getInstanceId(), null); } else { - sourceInstanceStatusData = worker().getFunctionAdmin().source().getSourceStatus( + sourceInstanceStatusData = worker().getFunctionAdmin().sources().getSourceStatus( assignment.getInstance().getFunctionMetaData().getFunctionDetails().getTenant(), assignment.getInstance().getFunctionMetaData().getFunctionDetails().getNamespace(), assignment.getInstance().getFunctionMetaData().getFunctionDetails().getName(), diff --git a/pulsar-io/file/src/main/java/org/apache/pulsar/io/file/utils/GZipFiles.java b/pulsar-io/file/src/main/java/org/apache/pulsar/io/file/utils/GZipFiles.java index 0a31e317fea51..f9e173f682bb9 100644 --- a/pulsar-io/file/src/main/java/org/apache/pulsar/io/file/utils/GZipFiles.java +++ b/pulsar-io/file/src/main/java/org/apache/pulsar/io/file/utils/GZipFiles.java @@ -44,20 +44,15 @@ public class GZipFiles { * Returns true if the given file is a gzip file. */ public static boolean isGzip(File f) { - - InputStream input = null; - try { - input = new FileInputStream(f); + try (InputStream input = new FileInputStream(f)) { PushbackInputStream pb = new PushbackInputStream(input, 2); - byte [] signature = new byte[2]; + byte[] signature = new byte[2]; int len = pb.read(signature); //read the signature pb.unread(signature, 0, len); //push back the signature to the stream // check if matches standard gzip magic number - return (signature[ 0 ] == (byte) 0x1f && signature[1] == (byte) 0x8b); + return (signature[0] == (byte) 0x1f && signature[1] == (byte) 0x8b); } catch (final Exception e) { return false; - } finally { - IOUtils.closeQuietly(input); } } diff --git a/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaAbstractSource.java b/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaAbstractSource.java index 490690a2b5f78..0a0074c9f81cf 100644 --- a/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaAbstractSource.java +++ b/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaAbstractSource.java @@ -39,6 +39,8 @@ import org.apache.pulsar.io.core.SourceContext; import org.slf4j.Logger; import org.slf4j.LoggerFactory; + +import java.time.Duration; import java.util.Objects; import java.util.Collections; import java.util.Map; @@ -152,7 +154,7 @@ public void start() { consumer.subscribe(Collections.singletonList(kafkaSourceConfig.getTopic())); LOG.info("Kafka source started."); while (running) { - ConsumerRecords consumerRecords = consumer.poll(1000); + ConsumerRecords consumerRecords = consumer.poll(Duration.ofSeconds(1L)); CompletableFuture[] futures = new CompletableFuture[consumerRecords.count()]; int index = 0; for (ConsumerRecord consumerRecord : consumerRecords) { diff --git a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarRecordCursor.java b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarRecordCursor.java index 52045d48abbce..218dc7c449e0e 100644 --- a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarRecordCursor.java +++ b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarRecordCursor.java @@ -66,11 +66,11 @@ import static java.util.concurrent.CompletableFuture.completedFuture; import static org.apache.pulsar.common.protocol.Commands.serializeMetadataAndPayload; +import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.eq; -import static org.mockito.Matchers.any; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/docker/ContainerExecResult.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/docker/ContainerExecResult.java index cbb5f9bfdd294..0e79bb6def058 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/docker/ContainerExecResult.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/docker/ContainerExecResult.java @@ -27,7 +27,7 @@ @Data(staticConstructor = "of") public class ContainerExecResult { - private final int exitCode; + private final long exitCode; private final String stdout; private final String stderr; diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/docker/ContainerExecResultBytes.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/docker/ContainerExecResultBytes.java index 0f29f187da601..402ad4e0ab254 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/docker/ContainerExecResultBytes.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/docker/ContainerExecResultBytes.java @@ -26,7 +26,7 @@ @Data(staticConstructor = "of") public class ContainerExecResultBytes { - private final int exitCode; + private final long exitCode; private final byte[] stdout; private final byte[] stderr; diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/KafkaSinkTester.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/KafkaSinkTester.java index 015897748a76b..ab66d3dead0e5 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/KafkaSinkTester.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/KafkaSinkTester.java @@ -22,6 +22,7 @@ import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertTrue; +import java.time.Duration; import java.util.Arrays; import java.util.Iterator; import java.util.Map; @@ -105,7 +106,7 @@ public void prepareSink() throws Exception { public void validateSinkResult(Map kvs) { Iterator> kvIter = kvs.entrySet().iterator(); while (kvIter.hasNext()) { - ConsumerRecords records = kafkaConsumer.poll(1000); + ConsumerRecords records = kafkaConsumer.poll(Duration.ofSeconds(1L)); log.info("Received {} records from kafka topic {}", records.count(), kafkaTopicName); if (records.isEmpty()) { diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/utils/DockerUtils.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/utils/DockerUtils.java index 9246c5802f3d1..2dbf77d19cd1a 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/utils/DockerUtils.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/utils/DockerUtils.java @@ -267,7 +267,7 @@ public void onComplete() { LOG.info("DOCKER.exec({}:{}): Done", containerName, cmdString); InspectExecResponse resp = waitForExecCmdToFinish(dockerClient, execId); - int retCode = resp.getExitCode(); + long retCode = resp.getExitCodeLong(); ContainerExecResult result = ContainerExecResult.of( retCode, stdout.toString(), @@ -342,7 +342,7 @@ public void onComplete() { future.join(); InspectExecResponse resp = waitForExecCmdToFinish(dockerClient, execId); - int retCode = resp.getExitCode(); + long retCode = resp.getExitCodeLong(); ContainerExecResultBytes result = ContainerExecResultBytes.of( retCode, @@ -356,9 +356,9 @@ public void onComplete() { return result; } - public static CompletableFuture runCommandAsyncWithLogging(DockerClient dockerClient, + public static CompletableFuture runCommandAsyncWithLogging(DockerClient dockerClient, String containerId, String... cmd) { - CompletableFuture future = new CompletableFuture<>(); + CompletableFuture future = new CompletableFuture<>(); String execId = dockerClient.execCreateCmd(containerId) .withCmd(cmd) .withAttachStderr(true) @@ -392,7 +392,7 @@ public void onError(Throwable throwable) { public void onComplete() { LOG.info("DOCKER.exec({}:{}): Done", containerName, cmdString); InspectExecResponse resp = waitForExecCmdToFinish(dockerClient, execId); - int retCode = resp.getExitCode(); + long retCode = resp.getExitCodeLong(); LOG.info("DOCKER.exec({}:{}): completed with {}", containerName, cmdString, retCode); future.complete(retCode); } diff --git a/tiered-storage/jcloud/src/main/java/org/apache/bookkeeper/mledger/offload/jcloud/provider/JCloudBlobStoreProvider.java b/tiered-storage/jcloud/src/main/java/org/apache/bookkeeper/mledger/offload/jcloud/provider/JCloudBlobStoreProvider.java index 49dabb261212c..341e92921ab3f 100644 --- a/tiered-storage/jcloud/src/main/java/org/apache/bookkeeper/mledger/offload/jcloud/provider/JCloudBlobStoreProvider.java +++ b/tiered-storage/jcloud/src/main/java/org/apache/bookkeeper/mledger/offload/jcloud/provider/JCloudBlobStoreProvider.java @@ -119,8 +119,9 @@ public BlobStore getBlobStore(TieredStorageConfiguration config) { public void buildCredentials(TieredStorageConfiguration config) { if (config.getCredentials() == null) { try { - String gcsKeyContent = Files.toString( - new File(config.getConfigProperty(GCS_ACCOUNT_KEY_FILE_FIELD)), Charset.defaultCharset()); + String gcsKeyContent = Files.asCharSource( + new File(config.getConfigProperty(GCS_ACCOUNT_KEY_FILE_FIELD)), + Charset.defaultCharset()).read(); config.setProviderCredentials(() -> new GoogleCredentialsFromJson(gcsKeyContent).get()); } catch (IOException ioe) { log.error("Cannot read GCS service account credentials file: {}", diff --git a/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/BlobStoreBackedInputStreamTest.java b/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/BlobStoreBackedInputStreamTest.java index 272ad124225af..618b544ebb7a9 100644 --- a/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/BlobStoreBackedInputStreamTest.java +++ b/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/BlobStoreBackedInputStreamTest.java @@ -32,10 +32,9 @@ import org.apache.bookkeeper.mledger.offload.jcloud.impl.BlobStoreBackedInputStreamImpl; import org.jclouds.blobstore.BlobStore; import org.jclouds.blobstore.domain.Blob; -import org.jclouds.blobstore.options.GetOptions; import org.jclouds.io.Payload; import org.jclouds.io.Payloads; -import org.mockito.Matchers; +import org.mockito.ArgumentMatchers; import org.mockito.Mockito; import org.testng.Assert; import org.testng.annotations.Test; @@ -218,7 +217,7 @@ public void testSeekWithinCurrent() throws Exception { } verify(spiedBlobStore, times(1)) - .getBlob(Mockito.eq(BUCKET), Mockito.eq(objectKey), Matchers.anyObject()); + .getBlob(Mockito.eq(BUCKET), Mockito.eq(objectKey), ArgumentMatchers.any()); } @Test