IP Library Granted Patent US 11,886,225
Granted Patent B2
US 11,886,225 · App. 17/568,828 · Granted Jan 30, 2024

Message processing method and apparatus in distributed system

Inventors: Xiaoqin Xie (Beijing, CN); Kun Li (Shenzhen, CN)
Assignee: HUAWEI CLOUD COMPUTING TECHNOLOGIES CO., LTD.
G06F9/546G06F11/3476
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 11,886,225
App. No.
17/568,828
Granted
Jan 30, 2024
Kind
B2
Abstract

In a message processing method, a message stream platform receives a plurality of log messages from a production platform, and the plurality of log messages are used to record information about a first service executed by the production platform. The message stream platform stores the plurality of log messages in a log file based on time segments. Then, the message stream platform may send messages in a same time segment in the plurality of log files to a consumption platform based on the time segment mark. According to the message processing method, the messages in the same time segment in the plurality of log files may be sent to the consumption platform to ensure that the consumption platform obtains the log messages generated in the same time segment.

Claims (70)

1. A method implemented by a message stream platform and comprising:

receiving, from a production platform, log messages that record information about a first service executed by the production platform;

allocating time segment marks to the log messages, wherein the time segment marks identify segments of time;

storing the log messages in a log file based on time segments and the time segment marks, wherein the time segments comprise a first time segment;

sending, to a consumption platform, the log messages that are related to the first service and that are in the first time segment;

receiving, from the production platform, an association request comprising a shard identifier of a second shard, wherein the second shard is based on splitting a first shard or aggregating the first shard and a third shard, and wherein the first shard and the third shard are a subset formed by the log messages during execution of the first service; and

setting, according to the association request, new time segment marks for subsequently received log messages related to the first service.

2. The method of claim 1 , further comprising further sending the log messages in response to a request from the consumption platform.

3. The method of claim 1 , wherein the time segments further comprise a second time segment, and wherein the method further comprises:

further sending, to the consumption platform based on a subscription relationship, the log messages that are related to the first service and that are in the first time segment; and

sending, to the consumption platform, the log messages that are related to the first service and that are in the second time segment.

4. The method of claim 1 , wherein allocating the time segment marks comprises:

periodically updating the time segment marks to obtain updated time segment marks; and

marking time segments for subsequently received log messages using the updated time segment marks.

5. The method of claim 4 , further comprising:

receiving, from the production platform, an association request comprising shard identifiers of both a first shard and a second shard or of both the second shard and a third shard, wherein the first shard and the third shard are a subset formed by the log messages during execution of the first service, and wherein the second shard is based on splitting the first shard or aggregating the first shard and the third shard; and

setting, according to the association request, new time segment marks for subsequently received log messages belonging to the first shard and the second shard.

6. A message stream platform comprising:

a memory configured to store instructions; and

one or more processors coupled to the memory and configured to execute the instructions to cause the message stream platform to:

receive, from a production platform, log messages that record information about a first service executed by the production platform;

allocate time segment marks to the log messages, wherein the time segment marks identify segments of time;

store the log messages in a log file based on time segments and the time segment marks, wherein the time segments comprise a first time segment;

send, to a consumption platform, the log messages that are related to the first service and that are in the first time segment;

receive, from the production platform, an association request comprising a shard identifier of a second shard, wherein the second shard is based on splitting a first shard or aggregating the first shard and a third shard, and wherein the first shard and the third shard are a subset formed by the log messages during execution of the first service; and

set, according to the association request, new time segment marks for subsequently received log messages related to the first service.

7. The message stream platform of claim 6 , wherein the one or more processors are further configured to execute the instructions to cause the message stream platform to further send the log messages in response to a request from the consumption platform.

8. The message stream platform of claim 6 , wherein the time segments further comprise a second time segment, and wherein the one or more processors are further configured to execute the instructions to cause the message stream platform to:

further send, to the consumption platform based on a subscription relationship, the log messages that are related to the first service and that are in the first time segment; and

send, to the consumption platform, the log messages that are related to the first service and that are in the second time segment.

9. The message stream platform of claim 6 , wherein the one or more processors are further configured to execute instructions to cause the message stream platform to:

periodically update the time segment marks to obtain updated time segment marks; and

mark time segments for subsequently received log messages using the updated time segment marks.

10. The message stream platform of claim 9 , wherein the one or more processors are further configured to execute the instructions to cause the message stream platform to:

receive, from the production platform, an association request comprising shard identifiers of both a first shard and a second shard or of both the second shard and a third shard, wherein the first shard and the third shard are a subset formed by the log messages about executing the first service, and wherein the second shard is based on splitting the first shard or aggregating the first shard and the third shard; and

set, according to the association request, new time segment marks for subsequently received log messages belonging to the first shard and the second shard.

11. A computer program product comprising instructions that are stored on a non-transitory computer-readable medium and that, when executed by one or more processors, cause a message stream platform to:

receive, from a production platform, log messages that record information about a first service executed by the production platform;

allocate time segment marks to the log messages, wherein the time segment marks identify segments of time;

store the log messages in a log file based on time segments and the time segment marks, wherein the time segments comprise a first time segment;

send, to a consumption platform, the log messages that are related to the first service and that are in the first time segment;

receive, from the production platform, an association request comprising a shard identifier of a second shard, wherein the second shard is based on splitting a first shard or aggregating the first shard and a third shard, and wherein the first shard and the third shard are a subset formed by the log messages during execution of the first service; and

set, according to the association request, new time segment marks for subsequently received log messages related to the first service.

12. The computer program product of claim 11 , wherein the instructions, when executed by the one or more processors, further cause the message stream platform to further send the log messages in response to a request from the consumption platform.

13. The computer program product of claim 11 , wherein the time segments further comprise a second time segment, and wherein the instructions, when executed by the one or more processors, further cause the message stream platform to:

further send, to the consumption platform based on a subscription relationship, the log messages that are related to the first service and that are in the first time segment; and

send, to the consumption platform, the log messages that are related to the first service and that are in the second time segment.

14. The computer program product of claim 11 , wherein the instructions, when executed by the one or more processors, further cause the message stream platform to allocate the time segment marks by:

periodically updating the time segment marks to obtain updated time segment marks; and

marking time segments for subsequently received log messages using the updated time segment marks.

15. The computer program product of claim 14 , wherein the instructions, when executed by the one or more processors, further cause the message stream platform to:

receive, from the production platform, an association request comprising shard identifiers of both a first shard and a second shard or of both the second shard and a third shard, wherein the first shard and the third shard are a subset formed by the log messages during execution of the first service, and wherein the second shard is based on splitting the first shard or aggregating the first shard and the third shard; and

set, according to the association request, new time segment marks for subsequently received log messages belonging to the first shard and the second shard.

16. A non-transitory computer-readable medium comprising instructions that, when executed by one or more processors, cause a message stream platform to:

receive, from a production platform, log messages that record information about a first service executed by the production platform;

allocate time segment marks to the log messages, wherein the time segment marks identify segments of time;

store the log messages in a log file based on time segments and the time segment marks, wherein the time segments comprise a first time segment;

send, to a consumption platform, the log messages that are related to the first service and that are in the first time segment;

receive, from the production platform, an association request comprising a shard identifier of a second shard, wherein the second shard is based on splitting a first shard or aggregating the first shard and a third shard, and wherein the first shard and the third shard are a subset formed by the log messages during execution of the first service; and

set, according to the association request, new time segment marks for subsequently received log messages related to the first service.

17. The non-transitory computer-readable medium of claim 16 , wherein the instructions, when executed by the one or more processors, further cause the message stream platform to further send the log messages in response to a request from the consumption platform.

18. The non-transitory computer-readable medium of claim 16 , wherein the time segments further comprise a second time segment, and wherein the instructions, when executed by the one or more processors, further cause the message stream platform to:

further send, to the consumption platform based on a subscription relationship, the log messages that are related to the first service and that are in the first time segment; and

send, to the consumption platform, the log messages that are related to the first service and that are in the second time segment.

19. The non-transitory computer-readable medium of claim 16 , wherein the instructions, when executed by the one or more processors, further cause the message stream platform to allocate the time segment marks by:

periodically updating the time segment marks to obtain updated time segment marks; and

marking time segments for subsequently received log messages using the updated time segment marks.

20. The non-transitory computer-readable medium of claim 19 , wherein the instructions, when executed by the one or more processors, further cause the message stream platform to:

receive, from the production platform, an association request comprising shard identifiers of both a first shard and a second shard or of both the second shard and a third shard, wherein the first shard and the third shard are a subset formed by the log messages during execution of the first service, and wherein the second shard is based on splitting the first shard or aggregating the first shard and the third shard; and

set, according to the association request, new time segment marks for subsequently received log messages belonging to the first shard and the second shard.

Assignments (2)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Mar 1, 2022
From: HUAWEI TECHNOLOGIES CO., LTD.
To: HUAWEI CLOUD COMPUTING TECHNOLOGIES CO., LTD.
Reel/Frame 059267/0088 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jan 5, 2022
From: XIE, XIAOQIN; LI, KUN
To: HUAWEI TECHNOLOGIES CO., LTD.
Reel/Frame 058553/0221 →