diff --git a/build.gradle b/build.gradle index dcafcbf687..d40deeed55 100644 --- a/build.gradle +++ b/build.gradle @@ -263,10 +263,10 @@ project(":samza-azure_$scalaSuffix") { apply plugin: 'java' dependencies { - compile "com.azure:azure-storage-blob:12.14.0" - compile "com.azure:azure-identity:1.3.6" - compile "com.microsoft.azure:azure-storage:5.3.1" - compile "com.microsoft.azure:azure-eventhubs:1.0.1" + compile "com.azure:azure-storage-blob:12.21.1" + compile "com.azure:azure-identity:1.8.1" + compile "com.microsoft.azure:azure-storage:8.6.6" + compile "com.microsoft.azure:azure-eventhubs:3.3.0" compile "com.fasterxml.jackson.core:jackson-core:$jacksonVersion" compile "io.dropwizard.metrics:metrics-core:3.1.2" compile "org.apache.avro:avro:$avroVersion" diff --git a/samza-azure/src/main/java/org/apache/samza/system/eventhub/EventHubConfig.java b/samza-azure/src/main/java/org/apache/samza/system/eventhub/EventHubConfig.java index e7e82aea13..1e21a21c28 100644 --- a/samza-azure/src/main/java/org/apache/samza/system/eventhub/EventHubConfig.java +++ b/samza-azure/src/main/java/org/apache/samza/system/eventhub/EventHubConfig.java @@ -187,7 +187,7 @@ public String getStreamEntityPath(String systemName, String streamName) { /** * Get the number of client threads, This is used to create the ThreadPool executor that is passed to the - * {@link EventHubClient#create} + * {@link EventHubClient#createFromConnectionStringSync} * @param systemName Name of the system. * @return Num of client threads to use. */ diff --git a/samza-azure/src/main/java/org/apache/samza/system/eventhub/SamzaEventHubClientManager.java b/samza-azure/src/main/java/org/apache/samza/system/eventhub/SamzaEventHubClientManager.java index ee515846b1..72289c1b55 100644 --- a/samza-azure/src/main/java/org/apache/samza/system/eventhub/SamzaEventHubClientManager.java +++ b/samza-azure/src/main/java/org/apache/samza/system/eventhub/SamzaEventHubClientManager.java @@ -26,7 +26,7 @@ import com.microsoft.azure.eventhubs.impl.RetryExponential; import com.microsoft.azure.eventhubs.RetryPolicy; import com.microsoft.azure.eventhubs.EventHubException; -import java.util.concurrent.ExecutorService; +import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.Executors; import org.apache.samza.SamzaException; import org.slf4j.Logger; @@ -55,7 +55,7 @@ public class SamzaEventHubClientManager implements EventHubClientManager { private final String sasKeyName; private final String sasKey; private final RetryPolicy retryPolicy; - private ExecutorService eventHubClientExecutor; + private ScheduledExecutorService eventHubClientExecutor; public SamzaEventHubClientManager(String eventHubNamespace, String entityPath, String sasKeyName, String sasKey, Integer numClientThreads) { @@ -85,8 +85,8 @@ public void init() { .setSasKey(sasKey); ThreadFactoryBuilder threadFactoryBuilder = new ThreadFactoryBuilder().setNameFormat("Samza EventHubClient Thread-%d").setDaemon(true); - eventHubClientExecutor = Executors.newFixedThreadPool(numClientThreads, threadFactoryBuilder.build()); - eventHubClient = EventHubClient.createSync(connectionStringBuilder.toString(), retryPolicy, eventHubClientExecutor); + eventHubClientExecutor = Executors.newScheduledThreadPool(numClientThreads, threadFactoryBuilder.build()); + eventHubClient = EventHubClient.createFromConnectionStringSync(connectionStringBuilder.toString(), retryPolicy, eventHubClientExecutor); } catch (IOException | EventHubException e) { String msg = String.format("Creation of EventHub client failed for eventHub EntityPath: %s on remote host %s:%d", entityPath, remoteHost, ClientConstants.AMQPS_PORT); diff --git a/samza-azure/src/main/java/org/apache/samza/system/eventhub/consumer/EventHubSystemConsumer.java b/samza-azure/src/main/java/org/apache/samza/system/eventhub/consumer/EventHubSystemConsumer.java index a6d975f5b7..91eab33eef 100644 --- a/samza-azure/src/main/java/org/apache/samza/system/eventhub/consumer/EventHubSystemConsumer.java +++ b/samza-azure/src/main/java/org/apache/samza/system/eventhub/consumer/EventHubSystemConsumer.java @@ -26,6 +26,7 @@ import com.microsoft.azure.eventhubs.EventPosition; import com.microsoft.azure.eventhubs.PartitionReceiveHandler; import com.microsoft.azure.eventhubs.PartitionReceiver; +import com.microsoft.azure.eventhubs.ReceiverOptions; import com.microsoft.azure.eventhubs.impl.ClientConstants; import java.time.Duration; import java.time.Instant; @@ -280,13 +281,13 @@ private synchronized void initializeEventHubsManagers() { } else { // EventHub will return the first message AFTER the offset that was specified in the fetch request. // If no such offset exists Eventhub will return an error. + ReceiverOptions receiverOptions = new ReceiverOptions(); + receiverOptions.setPrefetchCount(prefetchCount); receiver = eventHubClientManager.getEventHubClient() .createReceiver(consumerGroup, partitionId.toString(), EventPosition.fromOffset(offset, /* inclusiveFlag */false)).get(DEFAULT_EVENTHUB_CREATE_RECEIVER_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS); } - receiver.setPrefetchCount(prefetchCount); - PartitionReceiveHandler handler = new PartitionReceiverHandlerImpl(ssp, eventReadRates.get(streamId), eventByteReadRates.get(streamId), consumptionLagMs.get(streamId), readErrors.get(streamId), interceptors.getOrDefault(streamId, null), @@ -375,11 +376,11 @@ private void renewPartitionReceiver(SystemStreamPartition ssp) { streamPartitionReceivers.get(ssp).close().get(DEFAULT_SHUTDOWN_TIMEOUT_MILLIS, TimeUnit.MILLISECONDS); // Recreate receiver + ReceiverOptions receiverOptions = new ReceiverOptions(); + receiverOptions.setPrefetchCount(prefetchCount); PartitionReceiver receiver = eventHubClientManager.getEventHubClient() .createReceiverSync(consumerGroup, partitionId.toString(), - EventPosition.fromOffset(offset, !offset.equals(EventHubSystemConsumer.START_OF_STREAM))); - - receiver.setPrefetchCount(prefetchCount); + EventPosition.fromOffset(offset, !offset.equals(EventHubSystemConsumer.START_OF_STREAM)), receiverOptions); // Timeout for EventHubClient receive receiver.setReceiveTimeout(DEFAULT_EVENTHUB_RECEIVER_TIMEOUT); diff --git a/samza-azure/src/test/java/org/apache/samza/system/eventhub/MockEventData.java b/samza-azure/src/test/java/org/apache/samza/system/eventhub/MockEventData.java index 1e3d4f5cfd..1b57c10bb4 100644 --- a/samza-azure/src/test/java/org/apache/samza/system/eventhub/MockEventData.java +++ b/samza-azure/src/test/java/org/apache/samza/system/eventhub/MockEventData.java @@ -71,4 +71,34 @@ public Map getProperties() { public EventData.SystemProperties getSystemProperties() { return overridedSystemProperties; } + + @Override + public void setSystemProperties(SystemProperties props) { + this.overridedSystemProperties = props; + } + + @Override + public int compareTo(EventData that) { + return Long.compare( + this.getSystemProperties().getSequenceNumber(), + that.getSystemProperties().getSequenceNumber() + ); + } + + @Override + public boolean equals(Object o) { + if (this == o) { + return true; + } + if (!(o instanceof MockEventData)) { + return false; + } + MockEventData that = (MockEventData) o; + return Objects.equals(eventData, that.eventData); + } + + @Override + public int hashCode() { + return Objects.hash(eventData); + } } diff --git a/samza-tools/src/main/java/org/apache/samza/tools/EventHubConsoleConsumer.java b/samza-tools/src/main/java/org/apache/samza/tools/EventHubConsoleConsumer.java index a72c8b77f6..63628f71fc 100644 --- a/samza-tools/src/main/java/org/apache/samza/tools/EventHubConsoleConsumer.java +++ b/samza-tools/src/main/java/org/apache/samza/tools/EventHubConsoleConsumer.java @@ -102,7 +102,7 @@ private static void consumeEvents(String ehName, String namespace, String keyNam .setSasKeyName(keyName) .setSasKey(token); - EventHubClient client = EventHubClient.createSync(connStr.toString(), Executors.newFixedThreadPool(10)); + EventHubClient client = EventHubClient.createFromConnectionStringSync(connStr.toString(), Executors.newScheduledThreadPool(10)); EventHubRuntimeInformation runTimeInfo = client.getRuntimeInformation().get(); int numPartitions = runTimeInfo.getPartitionCount();