Menangani lonjakan sementara dengan kontrol alur

Pipeline data terkadang mengalami lonjakan traffic yang dipublikasikan. Lonjakan traffic dapat membebani pelanggan kecuali jika Anda siap menghadapinya. Solusi sederhana untuk menghindari lonjakan traffic adalah dengan meningkatkan resource pelanggan Pub/Sub secara dinamis untuk memproses lebih banyak pesan. Namun, solusi ini dapat meningkatkan biaya atau tidak berfungsi secara instan. Misalnya, Anda mungkin memerlukan banyak VM.

Kontrol alur di sisi pelanggan memungkinkan pelanggan mengatur kecepatan penyerapan pesan. Dengan demikian, kontrol alur menangani lonjakan traffic tanpa meningkatkan biaya atau hingga pelanggan di-scale up.

Kontrol alur adalah fitur yang tersedia di library klien tingkat tinggi Pub/Sub . Anda juga dapat menerapkan pemrograman kontrol alur Anda sendiri saat menggunakan library klien tingkat rendah.

Poin penting: Jika Anda mengalami lonjakan traffic yang tiba-tiba untuk pesan yang dipublikasikan, gunakan kontrol alur di klien pelanggan pull untuk mengatasi lonjakan sementara.

Kebutuhan akan kontrol alur menunjukkan bahwa pesan dipublikasikan dengan kecepatan yang lebih tinggi daripada yang digunakan. Jika skenario ini adalah status persisten, bukan lonjakan sementara dalam volume pesan, pertimbangkan untuk meningkatkan jumlah instance klien pelanggan.

Konfigurasi kontrol alur

Kontrol alur memungkinkan Anda mengonfigurasi jumlah maksimum byte yang dialokasikan untuk permintaan yang belum selesai, dan jumlah maksimum pesan yang belum selesai yang diizinkan. Tetapkan batas ini sesuai dengan kapasitas throughput mesin klien Anda.

Nilai default untuk variabel kontrol alur dan nama variabel mungkin berbeda di seluruh library klien. Misalnya, di library klien Java, variabel berikut mengonfigurasi kontrol alur:

  • setMaxOutstandingElementCount(). Menentukan jumlah maksimum pesan yang belum menerima konfirmasi atau konfirmasi negatif dari Pub/Sub.

  • setMaxOutstandingRequestBytes(). Menentukan ukuran maksimum pesan yang belum menerima konfirmasi atau konfirmasi negatif dari Pub/Sub.

Jika batas untuk setMaxOutstandingElementCount() atau setMaxOutstandingRequestBytes() terlampaui, klien pelanggan tidak akan menarik lebih banyak pesan. Perilaku ini akan terus berlanjut hingga pesan yang sudah ditarik dikonfirmasi atau dikonfirmasi negatif. Dengan demikian, kita dapat menyelaraskan throughput dengan biaya yang terkait dengan menjalankan lebih banyak pelanggan.

Contoh kode untuk kontrol alur

Untuk mengontrol kecepatan penerimaan pesan oleh klien pelanggan, gunakan fitur kontrol alur pelanggan. Fitur kontrol alur ini diilustrasikan dalam contoh berikut:

C++

Sebelum mencoba contoh ini, ikuti petunjuk penyiapan C++ di Panduan memulai: Menggunakan Library Klien. Untuk mengetahui informasi selengkapnya, lihat dokumentasi referensi API C++ Pub/Sub.

namespace pubsub = ::google::cloud::pubsub;
using ::google::cloud::future;
using ::google::cloud::Options;
using ::google::cloud::StatusOr;
auto sample = [](std::string project_id, std::string subscription_id) {
  // Change the flow control watermarks, by default the client library uses
  // 0 and 1,000 for the message count watermarks, and 0 and 10MiB for the
  // size watermarks. Recall that the library stops requesting messages if
  // any of the high watermarks are reached, and the library resumes
  // requesting messages when *both* low watermarks are reached.
  auto constexpr kMiB = 1024 * 1024L;
  auto subscriber = pubsub::Subscriber(pubsub::MakeSubscriberConnection(
      pubsub::Subscription(std::move(project_id), std::move(subscription_id)),
      Options{}
          .set<pubsub::MaxOutstandingMessagesOption>(1000)
          .set<pubsub::MaxOutstandingBytesOption>(8 * kMiB)));

  auto session = subscriber.Subscribe(
      [](pubsub::Message const& m, pubsub::AckHandler h) {
        std::move(h).ack();
        std::cout << "Received message " << m << "\n";
        PleaseIgnoreThisSimplifiesTestingTheSamples();
      });
  return std::make_pair(subscriber, std::move(session));
};

C#

Sebelum mencoba contoh ini, ikuti petunjuk penyiapan C# di Panduan memulai: Menggunakan Library Klien. Untuk mengetahui informasi selengkapnya, lihat dokumentasi referensi API C# Pub/Sub.


using Google.Api.Gax;
using Google.Cloud.PubSub.V1;
using System;
using System.Threading;
using System.Threading.Tasks;

public class PullMessagesWithFlowControlAsyncSample
{
    public async Task<int> PullMessagesWithFlowControlAsync(string projectId, string subscriptionId, bool acknowledge)
    {
        SubscriptionName subscriptionName = SubscriptionName.FromProjectSubscription(projectId, subscriptionId);
        int messageCount = 0;
        SubscriberClient subscriber = await new SubscriberClientBuilder
        {
            SubscriptionName = subscriptionName,
            Settings = new SubscriberClient.Settings
            {
                AckExtensionWindow = TimeSpan.FromSeconds(4),
                AckDeadline = TimeSpan.FromSeconds(10),
                FlowControlSettings = new FlowControlSettings(maxOutstandingElementCount: 100, maxOutstandingByteCount: 10240)
            }
        }.BuildAsync();
        // SubscriberClient runs your message handle function on multiple
        // threads to maximize throughput.
        Task startTask = subscriber.StartAsync((PubsubMessage message, CancellationToken cancel) =>
        {
            string text = message.Data.ToStringUtf8();
            Console.WriteLine($"Message {message.MessageId}: {text}");
            Interlocked.Increment(ref messageCount);
            return Task.FromResult(acknowledge ? SubscriberClient.Reply.Ack : SubscriberClient.Reply.Nack