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 @@ -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");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(() -> {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@
/**
* Default routing mode for messages to partition.
*
* <p>This logic is applied when the application is not setting a key {@link MessageBuilder#setKey(String)}
* <p>This logic is applied when the application is not setting a key {@link TypedMessageBuilder#key(String)}
* on a particular message.
*/
@InterfaceAudience.Public
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -121,7 +121,7 @@ public List<Metrics> 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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -152,7 +154,7 @@ public void start() {
consumer.subscribe(Collections.singletonList(kafkaSourceConfig.getTopic()));
LOG.info("Kafka source started.");
while (running) {
ConsumerRecords<Object, Object> consumerRecords = consumer.poll(1000);
ConsumerRecords<Object, Object> consumerRecords = consumer.poll(Duration.ofSeconds(1L));
CompletableFuture<?>[] futures = new CompletableFuture<?>[consumerRecords.count()];
int index = 0;
for (ConsumerRecord<Object, Object> consumerRecord : consumerRecords) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -105,7 +106,7 @@ public void prepareSink() throws Exception {
public void validateSinkResult(Map<String, String> kvs) {
Iterator<Map.Entry<String, String>> kvIter = kvs.entrySet().iterator();
while (kvIter.hasNext()) {
ConsumerRecords<String, String> records = kafkaConsumer.poll(1000);
ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofSeconds(1L));
log.info("Received {} records from kafka topic {}",
records.count(), kafkaTopicName);
if (records.isEmpty()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand Down Expand Up @@ -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,
Expand All @@ -356,9 +356,9 @@ public void onComplete() {
return result;
}

public static CompletableFuture<Integer> runCommandAsyncWithLogging(DockerClient dockerClient,
public static CompletableFuture<Long> runCommandAsyncWithLogging(DockerClient dockerClient,
String containerId, String... cmd) {
CompletableFuture<Integer> future = new CompletableFuture<>();
CompletableFuture<Long> future = new CompletableFuture<>();
String execId = dockerClient.execCreateCmd(containerId)
.withCmd(cmd)
.withAttachStderr(true)
Expand Down Expand Up @@ -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);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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: {}",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -218,7 +217,7 @@ public void testSeekWithinCurrent() throws Exception {
}

verify(spiedBlobStore, times(1))
.getBlob(Mockito.eq(BUCKET), Mockito.eq(objectKey), Matchers.<GetOptions>anyObject());
.getBlob(Mockito.eq(BUCKET), Mockito.eq(objectKey), ArgumentMatchers.any());
}

@Test
Expand Down