Implement a PubSub debounce mechanism

Viewed 267

I have different publishers publish to a PubSub Topic. Each message has a specific key. I would like to create subscribers that only pick up the latest message for each specific key within a defined interval. In other words, I would like to have some kind of debounce implemented for my subscribers.

Example (with debounce 2 seconds)

-(x)-(y)-(x)-------(z)-(z)---(x)----------------->                [Topic with messages]


  |-------|---------------|execute for x                          [Subscriber]
              2 seconds

      |---------------|execute for y                              [Subscriber]
          2 seconds      
             
                    |---|---------------|execute for z            [Subscriber]
                            2 seconds

                              |---------------|execute for x      [Subscriber]
                                   2 seconds

Ordered Execution Summary:
execute for message with key: y
execute for message with key: x
execute for message with key: z
execute for message with key: x

Implementation

// index.ts

import * as pubsub from '@google-cloud/pubsub';
import * as functions from 'firebase-functions';
import AbortController from 'node-abort-controller';

exports.Debouncer = functions
  .runWith({
    // runtimeOptions
  })
  .region('REGION')
  .pubsub.topic('TOPIC_NAME')
  .onPublish(async (message, context) => {
    const key = message.json.key;

    // when an equivalent topic is being received, cancel this calculation:
    const aborter = await abortHelper<any>(
      'TOPIC_NAME',
      (message) => message?.key === key
    ).catch((error) => {
      console.error('Failed to init abort helper', error);
      throw new Error('Failed to init abort helper');
    });

    await new Promise((resolve) => setTimeout(resolve, 2000));
    
    // here, run the EXECUTION for the key, unless an abortsignal from the abortHelper was received:
    // if(aborter.abortController.signal) ...

    aborter.teardown();

    /**
     * Subscribe to the first subscription found for the specified topic. Once a
     * message gets received that is matching `messageMatcher`, the returned
     * AbortController reflects the abortet state. Calling the returned teardown
     * will cancel the subscription.
     */
    async function abortHelper<TMessage>(
      topicName: string,
      messageMatcher: (message: TMessage) => boolean = () => true
    ) {
      const abortController = new AbortController();
      const pubSubClient = new pubsub.PubSub();
      const topic = pubSubClient.topic(topicName);
      const subscription = await topic
        .getSubscriptions()
        .then((subscriptionsResponse) => {
          // TODO use better approach to find or provide subscription
          const subscription = subscriptionsResponse?.[0]?.[0];
          if (!subscription) {
            throw new Error('no found subscription');
          }
          return subscription;
        });

      const listener = (message: TMessage) => {
        const matching = messageMatcher(message);
        if (matching) {
          abortController.abort();
          unsubscribeFromPubSubTopicSubscription();
        }
      };

      subscription.addListener('message', listener);

      return {
        teardown: () => {
          unsubscribeFromPubSubTopicSubscription();
        },
        abortController,
      };

      function unsubscribeFromPubSubTopicSubscription() {
        subscription.removeListener('message', listener);
      }
    }
  });


The initial idea was to register a cloud function to the topic. This cloud function itself then subscribes to the topic as well and waits for the defined interval. If it picks up a message with the same key during the interval, it exits the cloud function. Otherwise, it runs the execution.

Running inside the firebase-emulator this worked fine. However, on production random and hard to debug issues occurred most likely due to parallel execution of the functions.

What would be the best approach to implement such a system in a scalable way? (It does not necessarily have to be with PubSub.)

0 Answers
Related