[serve] Remove SLO code and blist dependency (#10075)

This commit is contained in:
Edward Oakes
2020-08-18 17:52:36 -05:00
committed by GitHub
parent 263df6163c
commit ba0f531da0
11 changed files with 8 additions and 230 deletions
+2 -21
View File
@@ -5,8 +5,6 @@ import time
from typing import DefaultDict, List
import pickle
import blist
from ray.exceptions import RayTaskError
import ray
@@ -24,7 +22,6 @@ class Query:
request_args,
request_kwargs,
request_context,
request_slo_ms,
call_method="__call__",
shard_key=None,
async_future=None,
@@ -36,10 +33,6 @@ class Query:
self.async_future = async_future
# Service level objective in milliseconds. This is expected to be the
# absolute time since unix epoch.
self.request_slo_ms = request_slo_ms
self.call_method = call_method
self.shard_key = shard_key
self.is_shadow_query = is_shadow_query
@@ -59,11 +52,6 @@ class Query:
kwargs = pickle.loads(value)
return Query(**kwargs)
# adding comparator fn for maintaining an
# ascending order sorted list w.r.t request_slo_ms
def __lt__(self, other):
return self.request_slo_ms < other.request_slo_ms
def _make_future_unwrapper(client_futures: List[asyncio.Future],
host_future: asyncio.Future):
@@ -111,7 +99,7 @@ class Router:
# backend_name -> worker replica tag queue
self.worker_queues: DefaultDict[deque[str]] = defaultdict(deque)
# backend_name -> worker payload queue
self.backend_queues = defaultdict(blist.sortedlist)
self.backend_queues = defaultdict(deque)
# -- Metadata -- #
@@ -184,18 +172,11 @@ class Router:
logger.debug("Received a request for endpoint {}".format(endpoint))
self.num_router_requests.labels(endpoint=endpoint).add()
# check if the slo specified is directly the
# wall clock time
if request_meta.absolute_slo_ms is not None:
request_slo_ms = request_meta.absolute_slo_ms
else:
request_slo_ms = request_meta.adjust_relative_slo_ms()
request_context = request_meta.request_context
query = Query(
request_args,
request_kwargs,
request_context,
request_slo_ms,
call_method=request_meta.call_method,
shard_key=request_meta.shard_key,
async_future=asyncio.get_event_loop().create_future())
@@ -366,7 +347,7 @@ class Router:
backend, curr_queries))
continue
request = buffer_queue.pop(0)
request = buffer_queue.pop()
self.queries_counter[backend][backend_replica_tag] += 1
future = asyncio.get_event_loop().create_task(
self._do_query(backend, backend_replica_tag, request))