Skip to content

feat(pubsub): make lease management RPCs concurrent - #10238

Merged
hongalex merged 12 commits into
googleapis:mainfrom
hongalex:pubsub-concurrent-lease-management-rpc
Jun 18, 2024
Merged

hongalex merged 12 commits into
googleapis:mainfrom
hongalex:pubsub-concurrent-lease-management-rpc

Conversation

@hongalex

@hongalex hongalex commented May 20, 2024 •

Copy link
Copy Markdown
Member

This PR makes the lease management RPCs concurrent rather than serial. The previous behavior takes a list of ackIDs that need to be sent out, takes the first n messages (where n is the max batch size), issues the RPC, and then repeats until ackID slice is empty.

This PR improves performance for large amounts of ackIDs by concurrently sending out RPCs in a goroutine, and waiting for them all to complete in a waitgroup. This means rewriting the splitRequestIDs function into a new makeBatches function, that does the splitting / batch making upfront, rather than per iteration.

Fixes #9727

@hongalex
hongalex requested a review from shollyman as a code owner May 20, 2024 21:48
@hongalex
hongalex requested review from a team May 20, 2024 21:48
@product-auto-label product-auto-label Bot added the api: pubsub Issues related to the Pub/Sub API. label May 20, 2024

@alvarowolfx alvarowolfx left a comment •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

overall looks good to me. One thing that got me a bit concerned is the logic on methods sendAck and sendModAck are almost the same and has duplicate code that can lead to it being out of sync sometimes. What differs is how to track telemetry, how to send the ack and how to retry.

I have a suggestion, but I don't have the full context to know if those ack methods can diverse more in the future, but we can have a common sendAckWithFunc method that contains the logic and accept functions for acking, retrying and recordTelemetry.

Would be something like this:

func (it *messageIterator) sendAck(m map[string]*AckResult) {
	it.sendAckWithFunc(m, func(ctx context.Context, subName string, ackIds []string) error {
		return it.subc.Acknowledge(ctx, &pb.AcknowledgeRequest{
			Subscription: it.subName,
			AckIds:       ackIds,
		})
	}, it.retryAcks, func(ctx context.Context, toSend []string) {
		recordStat(it.ctx, AckCount, int64(len(toSend)))
		addAcks(toSend)
	})
}

func (it *messageIterator) sendModAck(m map[string]*AckResult, deadline time.Duration, logOnInvalid bool) {
	deadlineSec := int32(deadline / time.Second)
	it.sendAckWithFunc(m, func(ctx context.Context, subName string, ackIds []string) error {
		return it.subc.ModifyAckDeadline(ctx, &pb.ModifyAckDeadlineRequest{
			Subscription:       it.subName,
			AckDeadlineSeconds: deadlineSec,
			AckIds:             ackIds,
		})
	}, func(toRetry map[string]*ipubsub.AckResult) {
		it.retryModAcks(toRetry, deadlineSec, logOnInvalid)
	}, func(ctx context.Context, toSend []string) {
		if deadline == 0 {
			recordStat(it.ctx, NackCount, int64(len(toSend)))
		} else {
			recordStat(it.ctx, ModAckCount, int64(len(toSend)))
		}
		addModAcks(toSend, deadlineSec)
	})
}

type ackFunc = func(ctx context.Context, subName string, ackIds []string) error
type ackRecordStat = func(ctx context.Context, toSend []string)
type retryAckFunc = func(toRetry map[string]*ipubsub.AckResult)

func (it *messageIterator) sendAckWithFunc(m map[string]*AckResult, ackFunc ackFunc, retryAckFunc retryAckFunc, ackRecordStat ackRecordStat) {
	ackIDs := make([]string, 0, len(m))
	for k := range m {
		ackIDs = append(ackIDs, k)
	}
	it.eoMu.RLock()
	exactlyOnceDelivery := it.enableExactlyOnceDelivery
	it.eoMu.RUnlock()
	batches := makeBatches(ackIDs, ackIDBatchSize)
	wg := sync.WaitGroup{}

	for _, batch := range batches {
		wg.Add(1)
		go func(toSend []string) {
			defer wg.Done()
			ackRecordStat(it.ctx, toSend)
			// Use context.Background() as the call's context, not it.ctx. We don't
			// want to cancel this RPC when the iterator is stopped.
			cctx, cancel2 := context.WithTimeout(context.Background(), 60*time.Second)
			defer cancel2()
			err := ackFunc(cctx, it.subName, toSend)
			if exactlyOnceDelivery {
				resultsByAckID := make(map[string]*AckResult)
				for _, ackID := range toSend {
					resultsByAckID[ackID] = m[ackID]
				}

				st, md := extractMetadata(err)
				_, toRetry := processResults(st, resultsByAckID, md)
				if len(toRetry) > 0 {
					// Retry modacks/nacks in a separate goroutine.
					go func() {
						retryAckFunc(toRetry)
					}()
				}
			}
		}(batch)
	}
}

@hongalex

Copy link
Copy Markdown
Member Author

I haven't been a huge fan of callback based systems since it makes jumping around in the IDE more difficult. With that said, the tradeoffs could be worth it in this case since there is a large number of overlap in code. I'm changing these functions a fair bit as part of the otel tracing change, but I think it should still be mostly compatible even when wrapped with functions.

I think that if we end up needing more than these set of function arguments, it would be even less ergonomic but as-is, this seems like a good change.

@hongalex
hongalex enabled auto-merge (squash) June 18, 2024 06:29
@hongalex
hongalex merged commit 426a8c2 into googleapis:main Jun 18, 2024
@hongalex
hongalex deleted the pubsub-concurrent-lease-management-rpc branch December 10, 2025 22:49
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

api: pubsub Issues related to the Pub/Sub API.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

pubsub: Many acks/modacks could cause other acks/modacks to be delayed

2 participants