This message was deleted.
# ask-for-help
s
This message was deleted.
n
Code:
Copy code
import numpy as np
from <http://bentoml.io|bentoml.io> import NumpyNdarray
from pyspark.sql.types import StructType, StructField, FloatType, StringType
from pyspark.sql.types import IntegerType
from pyspark.sql.types import ArrayType
import random
import time
import bentoml

bento = bentoml.get("image_clip:latest")

schema = StructType([
    StructField("image", StringType(), True),
])


batch_size = 10_000

data = [["<https://tr.rbxcdn.com/051082206a180e1625becfdb75900aa2/420/420/Model/Png>"] for _ in range(batch_size)]

df = spark.createDataFrame(data, schema)

start = time.perf_counter()
results_df = bentoml.batch.run_in_spark(bento = bento, df = df, spark = spark)
results_df.show()
end = time.perf_counter()
print("Total Time: {tot} ms".format(tot=round((end-start)*1_000, 3)))
Error:
Copy code
An exception was thrown from the Python worker. Please see the stack trace below.
Traceback (most recent call last):
  File "/usr/local/lib/python3.7/site-packages/bentoml/_internal/batch/spark.py", line 87, in process
    func_output = client.call(api_name, func_input)
  File "/usr/local/lib/python3.7/site-packages/bentoml/_internal/client/__init__.py", line 54, in call
    inp, _bentoml_api=self._svc.apis[bentoml_api_name], **kwargs
  File "/usr/local/lib/python3.7/site-packages/bentoml/_internal/client/__init__.py", line 133, in _sync_call
    return asyncio.run(self._call(inp, _bentoml_api=_bentoml_api, **kwargs))
  File "/usr/lib64/python3.7/asyncio/runners.py", line 50, in run
    loop.close()
  File "/usr/lib64/python3.7/asyncio/base_events.py", line 587, in run_until_complete
    return future.result()
  File "/usr/local/lib/python3.7/site-packages/bentoml/_internal/client/http.py", line 180, in _call
    fake_req._headers = headers  # type: ignore (request._headers is property)
  File "/usr/local/lib64/python3.7/site-packages/aiohttp/client.py", line 1141, in __aenter__
    self._resp = await self._coro
  File "/usr/local/lib64/python3.7/site-packages/aiohttp/client.py", line 671, in _request
    raise
  File "/usr/local/lib64/python3.7/site-packages/aiohttp/client_reqrep.py", line 914, in start
    self._continue = None
  File "/usr/local/lib64/python3.7/site-packages/aiohttp/helpers.py", line 721, in __exit__
    raise asyncio.TimeoutError from None
concurrent.futures._base.TimeoutError
c
Hi @Nikhil Patel - the default max batch size in Spark might be too large for some ML models
try setting the max records per batch parameter in spark can often help with this issue
e.g.:
Copy code
spark.conf.set("spark.sql.execution.arrow.maxRecordsPerBatch", "1000")
the default I believe is 10k - it can get quite slow sometimes depending on the model
The reason
api_server.traffic.timeout=9000000
doesn’t work is because you also need to set the runner time out to a larger value
Copy code
runners:
  traffic:
    timeout: ....
I’d recommend testing the optimal batch size for your model on the compute resources available on your spark nodes, and configure both api_server and runner timeout accordingly, based on the inference time per batch
n
Thanks for your response @Chaoyu! I've tried reducing the maxRecordsPerBatch to 10 and raised both timeouts by setting the BENTOML_CONFIG_OPTIONS variable to
"api_server.traffic.timeout=9000000 runners.traffic.timeout=9000000"
, but I still seem to timeout. This even happens when I run the toy iris classifier example on 320,000 data points. Is the value I'm assigning to BENTOML_CONFIG_OPTIONS formatted correctly?
c
It may has to do with the config not serialized correctly on worker nodes 🤔 @sauyon could you help take a look?
n
It looks like some people may have resolved this issue by setting
mb_max_latency
to a higher number, but I'm not sure how to do this since
Copy code
@svc.api(
    input=NumpyNdarray.from_sample(
        np.array([[4.9, 3.0, 1.4, 0.2]], dtype=np.double), enforce_shape=False),
    output=NumpyNdarray(dtype = np.double, shape = [1], 
    mb_max_latency=10000000, 
    mb_max_batch_size=10, 
    batch=True),
)
throws:
Copy code
TypeError: NumpyNdarray.__init__() got an unexpected keyword argument 'mb_max_latency'
(oops, sent the wrong code and error above. Assume that mb_max_latency etc. are passed into svc.api instead of NumpyNdarray())
s
Did you set that in the driver or worker node? We could probably support configuration propagating to workers, but for now maybe you can set
spark.executorEnv.BENTOML_CONFIG_OPTIONS
if you did it on the driver.
n
Yeah spark.executorEnv.BENTOML_CONFIG_OPTIONS was set when it failed
s
Ah, ok! Seems like something that requires investigation, then. Do you want to open an issue or should I?