'distributor' -> 'batcher' + slow down new ID checks
parent
70bad811bb
commit
4937f625d2
|
@ -27,6 +27,11 @@ with open(os.path.join(DIR, "batcher_config.yaml")) as f:
|
|||
|
||||
|
||||
class IDBatcher:
|
||||
"""
|
||||
Periodically gets the newest comment/post ID from reddit, compares it to the
|
||||
last fetched one, and adds all intermediate values into a queue.
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
self.lock = threading.Lock()
|
||||
self.ids_queue = deque()
|
||||
|
@ -43,7 +48,7 @@ class IDBatcher:
|
|||
|
||||
self.fetch_new_ids()
|
||||
|
||||
schedule.every(5).seconds.do(self.fetch_new_ids)
|
||||
schedule.every(30).seconds.do(self.fetch_new_ids)
|
||||
|
||||
thread = threading.Thread(target=self.start_scheduler)
|
||||
thread.start()
|
||||
|
|
|
@ -126,11 +126,15 @@ class RedditClient:
|
|||
|
||||
|
||||
class FetcherClient:
|
||||
"""
|
||||
Retrieves lists of IDs from the batcher, fetches them from reddit, and submits the JSON back.
|
||||
"""
|
||||
|
||||
def __init__(self, client_id, reddit_client):
|
||||
self.client_id = client_id
|
||||
self.reddit_client = reddit_client
|
||||
self.channel = grpc.insecure_channel(
|
||||
config["distributor_uri"],
|
||||
config["batcher_uri"],
|
||||
options=[
|
||||
("grpc.max_send_message_length", 100 * 10**6),
|
||||
("grpc.max_receive_message_length", 1 * 10**6),
|
||||
|
|
|
@ -10,7 +10,6 @@ service IDService {
|
|||
rpc SubmitBatch (SubmitRequest) returns (SubmitResponse) {}
|
||||
}
|
||||
|
||||
// The BatchRequest message contains the client id.
|
||||
message BatchRequest {
|
||||
string client_id = 1;
|
||||
}
|
||||
|
|
Loading…
Reference in New Issue