This message was deleted.
# ask-for-help
s
This message was deleted.
j
My service code looks like this:
Copy code
from __future__ import annotations

from collections.abc import Collection

import numpy as np
import duckdb
import pandas as pd
import xgboost as xgb
import numpy.typing as npt
from <http://bentoml.io|bentoml.io> import NumpyNdarray, PandasDataFrame

import bentoml

CONN = duckdb.connect(database=':memory:')

INPUT_SPEC = PandasDataFrame(
    orient='columns',
    columns=[
        'user_id',
        'item_id',
    ],
    dtype={
        'user_id': 'int64',
        'item_id': 'int64',
    },
    default_format='parquet',
)


class XGBoostRunnable(bentoml.Runnable):
    SUPPORTED_RESOURCES = ('cpu',)
    SUPPORTS_CPU_MULTI_THREADING = False

    def __init__(self, booster: xgb.Booster) -> None:
        self.booster = booster

    @bentoml.Runnable.method(batchable=True)
    def predict(self, inp: pd.DataFrame) -> npt.NDArray[np.float32]:
        features = get_features(inp, self.booster.feature_names)
        return np.asarray(self.booster.predict(xgb.DMatrix(features)), dtype=np.float32)


booster = xgb.Booster(model_file='model.ubj')
if (feature_names := booster.feature_names) is None:
    raise ValueError('Feature names are required.')

runner = bentoml.Runner(
    XGBoostRunnable,
    runnable_init_params={'booster': booster},
    # TODO: Update this when <https://github.com/bentoml/BentoML/pull/4110>
    # is merged and released.
    max_batch_size=1000 * len(feature_names),
    max_latency_ms=1000,
    embedded=True,
)

svc = bentoml.Service('item_scorer', runners=[runner])


@svc.on_startup
def on_startup(context: bentoml.Context) -> None:
    CONN.execute("""
        CREATE TABLE data AS SELECT * FROM 'data.parquet';
        CREATE UNIQUE INDEX user_id_idx ON data (user_id);
    """)


@svc.on_shutdown
def on_shutdown(context: bentoml.Context) -> None:
    CONN.close()


def get_features(
    inp: pd.DataFrame,
    columns: Collection[str] | None = None,
) -> pd.DataFrame:
    column_refs = [f'data."{col}"' for col in columns] if columns else ['data.*']
    return CONN.execute(f"""
        SELECT {', '.join(column_refs)}
        FROM inp
        INNER JOIN data USING (user_id)
    """).fetch_df()


@svc.api(input=INPUT_SPEC, output=NumpyNdarray(dtype=np.float32))
async def score_items(inp: pd.DataFrame) -> npt.NDArray[np.float32]:
    return await runner.predict.async_run(inp)
But if I call the service with this code:
Copy code
runs = 1000


def run() -> list[float]:
    with concurrent.futures.ThreadPoolExecutor(max_workers=10) as executor:

        def work(inp):
            client = bentoml.client.Client.from_url('localhost:3000')
            ts = time.perf_counter()
            client.call(
                'score_items',
                inp=inp,
            )
            return time.perf_counter() - ts

        reqs = (
            pd.DataFrame({
                'user_id': int(np.random.randint(low=0, high=10_000_000, size=1)),
                'item_id': np.random.randint(0, 100, size=100).tolist(),
            }).explode('item_id').astype(np.int64)
            for _ in range(runs)
        )

        times = list(executor.map(work, reqs))

    return times
t
That looks like it should work. Could you check out the /metrics endpoint, there should be "*_adaptive_batch_size_bucket" metrics when batching is enabled
👀 1
j
and put a print in the runner code like
print(inp.shape)
it always comes out as
(100,2)
Copy code
# HELP bentoml_api_server_request_duration_seconds API HTTP request duration in seconds
# TYPE bentoml_api_server_request_duration_seconds histogram
bentoml_api_server_request_duration_seconds_sum{endpoint="/favicon.ico",http_response_code="404",service_name="item_scorer",service_version="not available"} 0.033906666008988395
bentoml_api_server_request_duration_seconds_bucket{endpoint="/favicon.ico",http_response_code="404",le="0.005",service_name="item_scorer",service_version="not available"} 0.0
bentoml_api_server_request_duration_seconds_bucket{endpoint="/favicon.ico",http_response_code="404",le="0.01",service_name="item_scorer",service_version="not available"} 0.0
bentoml_api_server_request_duration_seconds_bucket{endpoint="/favicon.ico",http_response_code="404",le="0.025",service_name="item_scorer",service_version="not available"} 0.0
bentoml_api_server_request_duration_seconds_bucket{endpoint="/favicon.ico",http_response_code="404",le="0.05",service_name="item_scorer",service_version="not available"} 1.0
bentoml_api_server_request_duration_seconds_bucket{endpoint="/favicon.ico",http_response_code="404",le="0.075",service_name="item_scorer",service_version="not available"} 1.0
bentoml_api_server_request_duration_seconds_bucket{endpoint="/favicon.ico",http_response_code="404",le="0.1",service_name="item_scorer",service_version="not available"} 1.0
bentoml_api_server_request_duration_seconds_bucket{endpoint="/favicon.ico",http_response_code="404",le="0.25",service_name="item_scorer",service_version="not available"} 1.0
bentoml_api_server_request_duration_seconds_bucket{endpoint="/favicon.ico",http_response_code="404",le="0.5",service_name="item_scorer",service_version="not available"} 1.0
bentoml_api_server_request_duration_seconds_bucket{endpoint="/favicon.ico",http_response_code="404",le="0.75",service_name="item_scorer",service_version="not available"} 1.0
bentoml_api_server_request_duration_seconds_bucket{endpoint="/favicon.ico",http_response_code="404",le="1.0",service_name="item_scorer",service_version="not available"} 1.0
bentoml_api_server_request_duration_seconds_bucket{endpoint="/favicon.ico",http_response_code="404",le="2.5",service_name="item_scorer",service_version="not available"} 1.0
bentoml_api_server_request_duration_seconds_bucket{endpoint="/favicon.ico",http_response_code="404",le="5.0",service_name="item_scorer",service_version="not available"} 1.0
bentoml_api_server_request_duration_seconds_bucket{endpoint="/favicon.ico",http_response_code="404",le="7.5",service_name="item_scorer",service_version="not available"} 1.0
bentoml_api_server_request_duration_seconds_bucket{endpoint="/favicon.ico",http_response_code="404",le="10.0",service_name="item_scorer",service_version="not available"} 1.0
bentoml_api_server_request_duration_seconds_bucket{endpoint="/favicon.ico",http_response_code="404",le="+Inf",service_name="item_scorer",service_version="not available"} 1.0
bentoml_api_server_request_duration_seconds_count{endpoint="/favicon.ico",http_response_code="404",service_name="item_scorer",service_version="not available"} 1.0
# HELP bentoml_api_server_request_total Total number of HTTP requests
# TYPE bentoml_api_server_request_total counter
bentoml_api_server_request_total{endpoint="/favicon.ico",http_response_code="404",service_name="item_scorer",service_version="not available"} 1.0
# HELP bentoml_api_server_request_in_progress Total number of HTTP requests in progress now
# TYPE bentoml_api_server_request_in_progress gauge
bentoml_api_server_request_in_progress{endpoint="/favicon.ico",service_name="item_scorer",service_version="not available"} 0.0
Doesn't seem to be there 🤔
t
yea, that's strange
j
hey @Judah Rand does sending multiple rows of input at once works? just trying to weed up the possibility one by one
j
I am sending 100 rows of input at a time
I'm sending 100 rows at a time 1000 times in parallel (using 10 threads) but seeing only 100 rows in any given execution of
predict
in the Runner
j
which bentoml version is this?
j
main
Any ideas? 🤔
Have you been able to reproduce my issue? Would it be useful for me to come up with a smaller reproducible example?
j
Was on a call just now, before that i tried with the official xgboost example in BentoML (not a custom runnable) and batching works there. so yes! a smaller reproducible example will be great to help us find the issue
also when you run
bentoml serve --debug
you should be able to see logs like this
Copy code
2023-08-15T23:50:22+0800 [DEBUG] [runner:booster_tree:1] Dynamic batching cork released, batch size: 1 (trace=c91a593124d80bdd6572e416f2060a3d,span=d40fb90262fa658d,sampled=0,service.name=booster_tree)
j
Simple service:
Copy code
from __future__ import annotations

import numpy as np
import pandas as pd
import numpy.typing as npt
from <http://bentoml.io|bentoml.io> import NumpyNdarray, PandasDataFrame

import bentoml

COLUMNS = ['user_id', 'item_id']
INPUT_SPEC = PandasDataFrame(
    orient='columns',
    columns=COLUMNS,
    dtype={
        'user_id': 'int64',
        'item_id': 'int64',
    },
    default_format='parquet',
)


class CustomRunnable(bentoml.Runnable):
    SUPPORTED_RESOURCES = ('cpu',)
    SUPPORTS_CPU_MULTI_THREADING = False

    @bentoml.Runnable.method(batchable=True)
    def predict(self, inp: pd.DataFrame) -> npt.NDArray[np.float32]:
        print("batch_shape:", inp.shape)
        return np.ones((inp.shape[0], 1), dtype=np.float32)


runner = bentoml.Runner(
    CustomRunnable,
    # TODO: Update this when <https://github.com/bentoml/BentoML/pull/4110>
    # is merged and released.
    max_batch_size=1000 * len(COLUMNS),
    max_latency_ms=1000,
    embedded=True,
)

svc = bentoml.Service('svc', runners=[runner])


@svc.api(input=INPUT_SPEC, output=NumpyNdarray(dtype=np.float32))
async def score_items(inp: pd.DataFrame) -> npt.NDArray[np.float32]:
    return await runner.predict.async_run(inp)
Simple client:
Copy code
import time
import concurrent.futures

import numpy as np
import pandas as pd

import bentoml

runs = 1000


def run() -> list[float]:
    with concurrent.futures.ThreadPoolExecutor(max_workers=10) as executor:

        def work(inp):
            client = bentoml.client.Client.from_url('localhost:3000')
            ts = time.perf_counter()
            client.call(
                'score_items',
                inp=inp,
            )
            return time.perf_counter() - ts

        reqs = (
            pd.DataFrame({
                'user_id': int(np.random.randint(low=0, high=10_000_000, size=1)),
                'item_id': np.random.randint(0, 100, size=100).tolist(),
            }).explode('item_id').astype(np.int64)
            for _ in range(runs)
        )

        times = list(executor.map(work, reqs))

    return times


if __name__ == '__main__':
    run()
Copy code
2023-08-15T23:50:22+0800 [DEBUG] [runner:booster_tree:1] Dynamic batching cork released, batch size: 1 (trace=c91a593124d80bdd6572e416f2060a3d,span=d40fb90262fa658d,sampled=0,service.name=booster_tree)
Aha! Interesting. I do see these when I use
embedded=False
but not when I use
embedded=True
!
💡 2
I was using an embedded runner because the performance guide suggests it: https://docs.bentoml.org/en/latest/guides/performance.html Does batching not work for embedded runners?
My batch_size does seem to still always be 1 though 🤔 It would be good to have more control over the batching strategy.
t
ah yea.... this actually makes sense now... "embedded" means that the model is instantiated inside of the api worker. which means yes, there wouldn't be batching possible. batching happens when you have many inputs coming into your api workers and the runner is in a separate process where the inputs are "batched" as they are received from multiple api workers
You do have a little bit of control over it with these configuration parameters: https://docs.bentoml.com/en/latest/guides/batching.html#configuring-batching
but I would think that you should see batching happening anyway if you're sending fast enough...
j
I'm still struggling to get it to batch. Really, all I'm trying to check now is that the batching does work. I just don't seem to be able to actually trigger it!
No I lie! It works!
🦜
It was just
embedded=True
that was tripping me up!
j
i think one simple hack that you can do is, adding a
asyncio.sleep(1)
in the runner method, to fake increase the latency of the endpoint
j
i think one simple hack that you can do is, adding a
asyncio.sleep(1)
in the runner method, to fake increase the latency of the endpoint
That's what I ended up doing 🤣
j
or
time.sleep
for sync
the performance gain of an embedded runner is due to saving an additional network call, but the tradeoff is the batching will be disabled. but one should probably benchmark according to their use case to see what setup works best!
thanks for pointing that out, we should update our documentation on this too
j
the performance gain of an embedded runner is due to saving an additional network call, but the tradeoff is the batching will be disabled. but one should probably benchmark according to their use case to see what setup works best! (edited)
Makes a lot of sense and I should have looked into that sooner. There isn't much documented about
embedded
though. Thanks for all your help!
j
yeap! we need to update that for sure! thanks for helping us to spot it too!