Traiter plus de messages avec le contrôle de simultanéité

Le contrôle de simultanéité est une fonctionnalité disponible dans la bibliothèque cliente de haut niveau Pub/Sub. Vous pouvez également implémenter votre propre contrôle de simultanéité lorsque vous utilisez une bibliothèque de bas niveau.

La prise en charge du contrôle de simultanéité dépend du langage de programmation de la bibliothèque cliente. Pour les implémentations de langage qui prennent en charge les threads parallèles, comme C++, Go et Java, les bibliothèques clientes font un choix par défaut pour le nombre de threads.

Ce choix n'est peut-être pas optimal pour votre application. Par exemple, si votre application d'abonné ne suit pas le volume de messages entrants et n'est pas liée au processeur, vous devez augmenter le nombre de threads. Pour les opérations de traitement des messages gourmandes en ressources processeur, il peut être judicieux de réduire le nombre de threads.

Cette page explique le concept de contrôle de simultanéité et comment configurer la fonctionnalité pour vos clients abonnés. Pour configurer vos clients éditeurs pour le contrôle de simultanéité, consultez Contrôle de simultanéité.

Configurations du contrôle de simultanéité

Les valeurs par défaut des variables de contrôle de simultanéité et les noms des variables peuvent varier d'une bibliothèque cliente à l'autre. Pour en savoir plus, consultez la documentation de référence de l'API. Par exemple, dans la bibliothèque cliente Java, les méthodes permettant de configurer le contrôle de la concurrence sont setParallelPullCount(), setExecutorProvider(), setSystemExecutorProvider() et setChannelProvider().

  • setParallelPullCount() vous permet de déterminer le nombre de flux à ouvrir. Vous pouvez ouvrir d'autres flux si votre client abonné peut gérer plus de données que celles envoyées sur un seul flux, soit 10 Mbit/s.

  • setExecutorProvider() vous permet de personnaliser le fournisseur d'exécution utilisé pour le traitement des messages. Par exemple, vous pouvez remplacer le fournisseur d'exécuteur par un fournisseur qui renvoie un seul exécuteur partagé avec un nombre limité de threads sur plusieurs clients abonnés. Cette configuration permet de limiter le nombre de threads créés. Le nombre total de threads utilisés pour le contrôle de simultanéité dépend du provider d'exécution transmis dans la bibliothèque cliente et du nombre de pull parallèles.

  • setSystemExecutorProvider() vous permet de personnaliser le fournisseur d'exécuteur utilisé pour la gestion des baux. En règle générale, vous ne configurez pas cette valeur, sauf si vous souhaitez utiliser le même fournisseur d'exécuteur dans setExecutorProvider et setSystemExecutorProvider. Par exemple, vous pouvez utiliser le même fournisseur d'exécuteur si vous avez un certain nombre d'abonnements à faible débit. L'utilisation de la même valeur limite le nombre de threads dans le client.

  • setChannelProvider() vous permet de personnaliser le fournisseur de canaux utilisé pour ouvrir des connexions à Pub/Sub. En règle générale, vous ne configurez pas cette valeur, sauf si vous souhaitez utiliser le même canal sur plusieurs clients abonnés. Si vous réutilisez un canal sur un trop grand nombre de clients, des erreurs GOAWAY ou ENHANCE_YOUR_CALM peuvent se produire. Si ces erreurs s'affichent dans les journaux de votre application ou dans Cloud Logging, créez d'autres canaux.

Exemples de code pour le contrôle de la simultanéité

C++

Avant d'essayer cet exemple, suivez les instructions de configuration pour C++ dans le guide de démarrage rapide : Utiliser les bibliothèques clientes. Pour en savoir plus, consultez la documentation de référence de l'API Pub/Sub C++.

namespace pubsub = ::google::cloud::pubsub;
using ::google::cloud::future;
using ::google::cloud::GrpcBackgroundThreadPoolSizeOption;
using ::google::cloud::Options;
using ::google::cloud::StatusOr;
auto sample = [](std::string project_id, std::string subscription_id) {
  // Create a subscriber with 16 threads handling I/O work, by default the
  // library creates `std::thread::hardware_concurrency()` threads.
  auto subscriber = pubsub::Subscriber(pubsub::MakeSubscriberConnection(
      pubsub::Subscription(std::move(project_id), std::move(subscription_id)),
      Options{}
          .set<pubsub::MaxConcurrencyOption>(8)
          .set<GrpcBackgroundThreadPoolSizeOption>(16)));

  // Create a subscription where up to 8 messages are handled concurrently. By
  // default the library uses `std::thread::hardware_concurrency()` as the
  // maximum number of concurrent callbacks.
  auto session = subscriber.Subscribe(
      [](pubsub::Message const& m, pubsub::AckHandler h) {
        // This handler executes in the I/O threads, applications could use,
        // std::async(), a thread-pool, or any other mechanism to transfer the
        // execution to other threads.
        std::cout << "Received message " << m << "\n";
        std::move(h).ack();
        PleaseIgnoreThisSimplifiesTestingTheSamples();
      });
  return std::make_pair(subscriber, std::move(session));
};

Go

L'exemple suivant utilise la version majeure de la bibliothèque cliente Go Pub/Sub (v2). Si vous utilisez toujours la bibliothèque v1, consultez le guide de migration vers la v2. Pour obtenir la liste des exemples de code v1, consultez les exemples de code obsolètes.

Avant d'essayer cet exemple, suivez les instructions de configuration pour Go dans le guide de démarrage rapide : Utiliser les bibliothèques clientes. Pour en savoir plus, consultez la documentation de référence de l'API Pub/Sub en langage Go.

import (
	"context"
	"fmt"
	"io"
	"sync/atomic"
	"time"

	"cloud.google.com/go/pubsub/v2"
)

func pullMsgsConcurrencyControl(w io.Writer, projectID, subID string) error {
	// projectID := "my-project-id"
	// subID := "my-sub"
	ctx := context.Background()
	client, err := pubsub.NewClient(ctx, projectID)
	if err != nil {
		return fmt.Errorf("pubsub.NewClient: %w", err)
	}
	defer client.Close()

	// client.Subscriber can be passed a subscription ID (e.g. "my-sub") or
	// a fully qualified name (e.g. "projects/my-project/subscriptions/my-sub").
	// If a subscription ID is provided, the project ID from the client is used.
	sub := client.Subscriber(subID)
	// NumGoroutines determines the number of streams sub.Receive will spawn to pull
	// messages. It is recommended to set this to 1, unless your throughput
	// is greater than 10 MB/s, as even having 1 stream can still result in
	// messages being handled asynchronously.
	sub.ReceiveSettings.NumGoroutines = 1
	// MaxOutstandingMessages limits the number of concurrent handlers of messages.
	// In this case, up to 8 unacked messages can be handled concurrently.
	sub.ReceiveSettings.MaxOutstandingMessages = 8

	// Receive messages for 10 seconds, which simplifies testing.
	// Comment this out in production, since `Receive` should
	// be used as a long running operation.
	ctx, cancel := context.WithTimeout(ctx, 10*time.Second)
	defer cancel()

	var received int32

	// Receive blocks until the context is cancelled or an error occurs.
	err = sub.Receive(ctx, func(_ context.Context, msg *pubsub.Message) {
		atomic.AddInt32(&received, 1)
		msg.Ack()
	})
	if err != nil {
		return fmt.Errorf("sub.Receive returned error: %w", err)
	}
	fmt.Fprintf(w, "Received %d messages\n", received)

	return nil
}

Java

Avant d'essayer cet exemple, suivez les instructions d'installation dans le langage Java se trouvant sur la page Démarrage rapide : utiliser des bibliothèques clientes. Pour en savoir plus, consultez la documentation de référence de l'API Pub/Sub en langage Java.


import com.google.api.gax.core.ExecutorProvider;
import com.google.api.gax.core.InstantiatingExecutorProvider;
import com.google.cloud.pubsub.v1.AckReplyConsumer;
import com.google.cloud.pubsub.v1.MessageReceiver;
import com.google.cloud.pubsub.v1.Subscriber;
import com.google.pubsub.v1.ProjectSubscriptionName;
import com.google.pubsub.v1.PubsubMessage;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;

public class SubscribeWithConcurrencyControlExample {
  public static void main(String... args) throws Exception {
    // TODO(developer): Replace these variables before running the sample.
    String projectId = "your-project-id";
    String subscriptionId = "your-subscription-id";

    subscribeWithConcurrencyControlExample(projectId, subscriptionId);
  }

  public static void subscribeWithConcurrencyControlExample(
      String projectId, String subscriptionId) {
    ProjectSubscriptionName subscriptionName =
        ProjectSubscriptionName.of(projectId, subscriptionId);

    // Instantiate an asynchronous message receiver.
    MessageReceiver receiver =
        (PubsubMessage message, AckReplyConsumer consumer) -> {
          // Handle incoming message, then ack the received message.
          System.out.println("Id: " + message.getMessageId());
          System.out.println("Data: " + message.getData().toStringUtf8());
          consumer.ack();
        };

    Subscriber subscriber = null;
    try {
      // Provides an executor service for processing messages. The default `executorProvider` used
      // by the subscriber has a default thread count of 5.
      ExecutorProvider executorProvider =
          InstantiatingExecutorProvider.newBuilder().setExecutorThreadCount(4).build();

      // `setParallelPullCount` determines how many StreamingPull streams the subscriber will open
      // to receive message. It defaults to 1. `setExecutorProvider` configures an executor for the
      // subscriber to process messages. Here, the subscriber is configured to open 2 streams for
      // receiving messages, each stream creates a new executor with 4 threads to help process the
      // message callbacks. In total 2x4=8 threads are used for message processing.
      subscriber =
          Subscriber.newBuilder(subscriptionName, receiver)
              .setParallelPullCount(2)
              .setExecutorProvider(executorProvider)
              .build();

      // Start the subscriber.
      subscriber.startAsync().awaitRunning();
      System.out.printf("Listening for messages on %s:\n", subscriptionName.toString());
      // Allow the subscriber to run for 30s unless an unrecoverable error occurs.
      subscriber.awaitTerminated(30, TimeUnit.SECONDS);
    } catch (TimeoutException timeoutException) {
      // Shut down the subscriber after 30s. Stop receiving messages.
      subscriber.stopAsync();
    }
  }
}

Ruby

L'exemple suivant utilise la bibliothèque cliente Ruby Pub/Sub v3. Si vous utilisez toujours la bibliothèque v2, consultez le guide de migration vers la v3. Pour obtenir la liste des exemples de code Ruby v2, consultez les exemples de code obsolètes.

Avant d'essayer cet exemple, suivez les instructions de configuration pour Ruby dans le guide de démarrage rapide : Utiliser les bibliothèques clientes. Pour en savoir plus, consultez la documentation de référence de l'API Pub/Sub en langage Ruby.

# subscription_id = "your-subscription-id"

pubsub = Google::Cloud::PubSub.new
subscriber = pubsub.subscriber subscription_id

# Use 2 threads for streaming, 4 threads for executing callbacks and 2 threads
# for sending acknowledgements and/or delays
listener = subscriber.listen streams: 2, threads: {
  callback: 4,
  push:     2
} do |received_message|
  puts "Received message: #{received_message.data}"
  received_message.acknowledge!
end

listener.start
# Let the main thread sleep for 60 seconds so the thread for listening
# messages does not quit
sleep 60
listener.stop.wait!

Rust

Avant d'essayer cet exemple, suivez les instructions de configuration pour Rust dans le guide de démarrage rapide : Utiliser les bibliothèques clientes. Pour en savoir plus, consultez la documentation de référence de l'API Pub/Sub Rust.

use google_cloud_pubsub::client::Subscriber;
use google_cloud_pubsub::subscriber::MessageStream;
use std::time::Duration;

pub async fn sample(project_id: &str, subscription_id: &str) -> anyhow::Result<()> {
    // Assume an available concurrency of 2 for this sample.
    const NCPU: usize = 2;

    let subscription_name = format!("projects/{project_id}/subscriptions/{subscription_id}");
    let client = Subscriber::builder()
        // Configure the subscriber to use multiple gRPC channels. This lets the
        // client multiplex its open streams and its acknowledgement RPCs.
        .with_grpc_subchannel_count(NCPU)
        .build()
        .await?;
    let tasks: Vec<_> = (0..2 * NCPU)
        .map(|i| {
            // Pub/Sub caps the throughput of a single stream to 10 MB/s. To achieve
            // higher throughput, you should open multiple streams.
            let stream = client.subscribe(&subscription_name).build();
            tokio::spawn(subscribe_task(i, stream))
        })
        .collect();

    for t in tasks {
        t.await??;
    }
    println!("done listening for messages");
    Ok(())
}

async fn subscribe_task(index: usize, mut stream: MessageStream) -> anyhow::Result<()> {
    // Terminate the example after 10 seconds.
    let shutdown_token = stream.shutdown_token();
    let shutdown = tokio::spawn(async move {
        tokio::time::sleep(Duration::from_secs(10)).await;
        shutdown_token.shutdown().await;
    });

    println!("listening for messages on stream {index}...");

    while let Some((m, h)) = stream.next().await.transpose()? {
        println!("received message: {m:?}");
        h.ack();
    }
    shutdown.await?;

    Ok(())
}

Étapes suivantes

Découvrez les autres options de distribution que vous pouvez configurer pour un abonnement :