IP Library › Granted Patent US 10,979,363
Granted Patent B2
US 10,979,363 · App. 16/801,853 · Granted Apr 13, 2021

Live resegmenting of partitions in distributed stream-processing platforms

Inventors: Andrey Efimov (Monroe, WA); John Christopher Petry (Seattle, WA); Julien Nicolas Dollon (Bothell, WA); Nathaniel Martin Glass (Bellevue, WA)
Assignee: Oracle International Corporation
H04L47/76H04L47/783H04L67/10H04L67/1095
View Patent ↗
Loading inventors, assignments & file history…
Monitor This Case
Get email alerts when status or documents change.
Order Certified Copies
Most orders are placed with the USPTO same day — all within 24 business hours.
Order via The Patent Place →
Pre-filled with this patent's details
Quick Facts
Patent No.
US 10,979,363
App. No.
16/801,853
Granted
Apr 13, 2021
Kind
B2
Abstract

Techniques for resegmenting a partition in a distributed stream-processing platform are provided. The techniques include receiving a trigger to move a partition of the distributed stream-processing platform from a first broker on a first set of physical resources to a second broker on a second a set of physical resources. In response to the trigger, the partition is allocated on the second broker, and the first broker is configured to redirect, to the second broker, requests for new messages after a last offset in the partition without replicating older messages before the last offset to the second broker. Idempotent produce metadata for the partition from the first broker is then merged into the second broker. Finally, metadata for processing requests for the partition is updated to include the second broker.

Claims (65)

1. A non-transitory computer readable medium comprising instructions which, when executed by one or more hardware processors, causes performance of operations comprising:

receiving metadata comprising:

a new broker for a partition in a distributed stream-processing platform; and

a last offset of messages in the partition that reside on an old broker for the partition;

storing, by a node in an interface layer of the distributed stream-processing platform, a mapping of the partition to the metadata;

directing, by the node to the old broker based on the mapping, a first read request for one or more messages in the partition before the last offset on the old broker; and

directing, by the node to the new broker based on the mapping, a second read request for one or more additional messages after the last offset on the old broker.

2. The medium of claim 1 , wherein the operations further comprise:

directing, by the node to the new broker based on the mapping, a write request for the partition.

3. The medium of claim 1 , wherein the operations further comprise:

upon receiving a trigger to move the partition from the old broker to the new broker, allocating the partition on the new broker; and

configuring the old broker to redirect, to the new broker, requests for new messages after the last offset in the partition without replicating older messages before the last offset to the new broker.

4. The medium of claim 3 , wherein the operations further comprise:

configuring the old broker to process read requests for old messages before the last offset in the partition during a retention period for the old messages.

5. The medium of claim 3 , wherein the operations further comprise:

after the old broker is configured to redirect the requests for the new messages after the last offset in the partition to the new broker, using idempotent produce metadata for a set of producers of the partition to validate, by the old broker, write requests for the partition; and

merging the idempotent produce metadata for the partition from the old broker into the new broker.

6. The medium of claim 5 , wherein the idempotent produce metadata comprises:

a producer identifier for a producer; and

a latest sequence number for the producer.

7. The medium of claim 3 , wherein configuring the old broker to redirect the requests for the new messages after the last offset in the partition to the new broker comprises:

updating the old broker with a redirect state and a redirect destination representing the new broker.

8. The medium of claim 3 , wherein redirecting the requests for the new messages after the last offset in the partition to the new broker comprises:

transmitting, with a request redirected to the new broker, the last offset in the partition for use in setting a base offset for the partition in the new broker.

9. The medium of claim 3 , wherein the trigger is received in response to a change in load on the old broker.

10. The medium of claim 1 , wherein receiving the metadata comprises:

receiving the metadata via a stream in the distributed stream-processing platform.

11. The medium of claim 1 , wherein storing the mapping of the partition to the metadata comprises:

storing a cluster containing the new broker in the mapping.

12. The medium of claim 1 , wherein storing the mapping of the partition to the metadata comprises:

storing a topic containing the partition in the mapping.

13. A method, comprising:

receiving metadata comprising:

a new broker for a partition in a distributed stream-processing platform; and

a last offset of messages in the partition that reside on an old broker for the partition;

storing, by a node in an interface layer of the distributed stream-processing platform, a mapping of the partition to the metadata;

directing, by the node to the old broker based on the mapping, a first read request for one or more messages in the partition before the last offset on the old broker; and

directing, by the node to the new broker based on the mapping, a second read request for one or more additional messages after the last offset on the old broker.

14. The method of claim 13 , further comprising:

directing, by the node to the new broker based on the mapping, write requests for the partition.

15. The method of claim 13 , further comprising:

upon receiving a trigger to move the partition from the old broker to the new broker, allocating the partition on the new broker;

configuring the old broker to redirect, to the new broker, requests for new messages after the last offset in the partition without replicating older messages before the last offset to the new broker; and

configuring the old broker to process read requests for old messages before the last offset in the partition during a retention period for the old messages.

16. The method of claim 13 , further comprising:

after the old broker is configured to redirect the requests for the new messages after the last offset in the partition to the new broker, using idempotent produce metadata for a set of producers of the partition to validate, by the old broker, write requests for the partition; and

merging the idempotent produce metadata for the partition from the old broker into the new broker.

17. The method of claim 16 , wherein the idempotent produce metadata comprises:

a producer identifier for a producer; and

a latest sequence number for the producer.

18. The method of claim 13 , wherein the trigger is received in response to a change in load on the old broker.

19. The method of claim 13 , wherein storing the mapping of the partition to the metadata comprises:

storing a cluster containing the new broker and a topic containing the partition in the mapping.

20. An apparatus, comprising:

one or more processors; and

memory storing instructions that, when executed by the one or more processors, cause the apparatus to:

receive metadata comprising:

a new broker for a partition in a distributed stream-processing platform;

a cluster containing the new broker;

a topic containing the partition; and

a last offset of messages in the partition that reside on an old broker for the partition;

store, by a node in an interface layer of the distributed stream-processing platform, a mapping of the partition to the metadata;

direct, by the node to the old broker based on the mapping, a first read request for one or more messages in the partition before the last offset on the old broker;

direct, by the node to the new broker based on the mapping, a second read request for one or more additional messages after the last offset on the old broker; and

direct, by the node to the new broker based on the mapping, write requests for the partition.

Continuity (3)
Continuation 15908465 · Feb 28, 2018
Provisional Application 62566370 · Sep 30, 2017
Related Publication 20200195572A1 · Jun 18, 2020
Cited By (1)
US 12,524,442