IP Library › Granted Patent US 10,608,951
Granted Patent B2
US 10,608,951 · App. 15/908,465 · Granted Mar 31, 2020

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,608,951
App. No.
15/908,465
Granted
Mar 31, 2020
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 (49)

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

upon receiving a trigger to move a partition of a 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, allocating the partition on the second broker;

configuring the first broker 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; and

updating metadata for processing requests for the partition to include the second broker.

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

merging idempotent produce metadata for the partition from the first broker into the second broker after the first broker is configured to redirect the requests for the new messages after the last offset in the partition to the second broker.

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

using the idempotent produce metadata to validate, at the first broker, write requests for the partition prior to merging the idempotent produce data from the first broker into the second broker.

4. The medium of claim 2 , wherein the idempotent produce metadata comprises:

a producer identifier for a producer; and

a latest sequence number for the producer.

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

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

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

using the metadata to direct read and write requests for the partition to the first and second brokers.

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

updating the first broker with a redirect state and a redirect destination representing the second broker.

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

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

9. The medium of claim 1 , wherein updating the metadata for processing requests for the partition to include the second broker comprises:

using a stream in the distributed stream-processing platform to propagate the metadata to a set of interface nodes in the distributed stream-processing platform.

10. The medium of claim 1 , wherein the trigger is received in response to a change in load on the first broker.

11. The medium of claim 1 , wherein the metadata further comprises a cluster, a topic, and the last offset.

12. A method, comprising:

upon receiving a trigger to move a partition of a 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, allocating the partition on the second broker;

configuring, by a computer system, the first broker 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; and

updating metadata for processing requests for the partition to include the second broker.

13. The method of claim 12 , further comprising:

merging idempotent produce metadata for the partition from the first broker into the second broker after the first broker is configured to redirect the requests for the new messages after the last offset in the partition to the second broker; and

using the idempotent produce metadata to validate, at the first broker, write requests for the partition prior to merging the idempotent produce data from the first broker into the second broker.

14. The method of claim 13 , wherein the idempotent produce metadata comprises:

a producer identifier for a producer; and

a latest sequence number for the producer.

15. The method of claim 12 , further comprising:

configuring the first 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 12 , further comprising:

using the metadata to direct read and write requests for the partition to the first and second brokers.

17. The method of claim 12 , wherein configuring the first broker to redirect the requests for the new messages after the last offset in the partition to the second broker comprises:

updating the first broker with a redirect state and a redirect destination representing the second broker.

18. The method of claim 12 , wherein redirecting the requests for the new messages after the last offset in the partition to the second broker comprises:

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

19. An apparatus, comprising:

one or more processors; and

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

upon receiving a trigger to move a partition of a 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, allocate the partition on the second broker;

configure the first broker 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; and

update metadata for processing requests for the partition to include the second broker.

20. The apparatus of claim 19 , wherein the memory further stores instructions that, when executed by the one or more processors, cause the apparatus to:

merge idempotent produce metadata for the partition from the first broker into the second broker after the first broker is configured to redirect the requests for the new messages after the last offset in the partition to the second broker.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Mar 2, 2018
From: EFIMOV, ANDREY; PETRY, JOHN CHRISTOPHER; DOLLON, JULIEN NICOLAS; GLASS, NATHANIEL MARTIN
To: ORACLE INTERNATIONAL CORPORATION
Reel/Frame 045094/0676 →
Continuity (2)
Provisional Application 62566370 · Sep 30, 2017
Related Publication 20190104082A1 · Apr 4, 2019