Skip to content
GitLab
Menu
Projects
Groups
Snippets
Loading...
Help
Help
Support
Community forum
Keyboard shortcuts
?
Submit feedback
Contribute to GitLab
Sign in / Register
Toggle navigation
Menu
Open sidebar
OpenDAS
dynamo
Commits
c30c6990
Unverified
Commit
c30c6990
authored
Apr 16, 2025
by
ishandhanani
Committed by
GitHub
Apr 17, 2025
Browse files
fix: direct clients vs dependancies (#704)
Co-authored-by:
Ziqi Fan
<
ziqif@nvidia.com
>
parent
f4780e85
Changes
2
Hide whitespace changes
Inline
Side-by-side
Showing
2 changed files
with
39 additions
and
16 deletions
+39
-16
examples/llm/components/processor.py
examples/llm/components/processor.py
+19
-8
examples/tensorrt_llm/components/processor.py
examples/tensorrt_llm/components/processor.py
+20
-8
No files found.
examples/llm/components/processor.py
View file @
c30c6990
...
@@ -67,6 +67,7 @@ class Processor(ProcessMixIn):
...
@@ -67,6 +67,7 @@ class Processor(ProcessMixIn):
self
.
tokenizer
,
self
.
model_config
self
.
tokenizer
,
self
.
model_config
)
)
self
.
min_workers
=
1
self
.
min_workers
=
1
print
(
f
"Processor init:
{
self
.
engine_args
.
router
}
"
)
def
_create_tokenizer
(
self
,
engine_args
:
AsyncEngineArgs
)
->
AnyTokenizer
:
def
_create_tokenizer
(
self
,
engine_args
:
AsyncEngineArgs
)
->
AnyTokenizer
:
"""Create a TokenizerGroup using engine arguments similar to VLLM's approach"""
"""Create a TokenizerGroup using engine arguments similar to VLLM's approach"""
...
@@ -93,6 +94,15 @@ class Processor(ProcessMixIn):
...
@@ -93,6 +94,15 @@ class Processor(ProcessMixIn):
.
client
()
.
client
()
)
)
if
self
.
engine_args
.
router
==
"kv"
:
router_ns
,
router_name
=
Router
.
dynamo_address
()
# type: ignore
self
.
router_client
=
(
await
runtime
.
namespace
(
router_ns
)
.
component
(
router_name
)
.
endpoint
(
"generate"
)
.
client
()
)
await
check_required_workers
(
self
.
worker_client
,
self
.
min_workers
)
await
check_required_workers
(
self
.
worker_client
,
self
.
min_workers
)
self
.
etcd_kv_cache
=
await
EtcdKvCache
.
create
(
self
.
etcd_kv_cache
=
await
EtcdKvCache
.
create
(
...
@@ -117,15 +127,16 @@ class Processor(ProcessMixIn):
...
@@ -117,15 +127,16 @@ class Processor(ProcessMixIn):
)
=
await
self
.
_parse_raw_request
(
raw_request
)
)
=
await
self
.
_parse_raw_request
(
raw_request
)
router_mode
=
(
await
self
.
etcd_kv_cache
.
get
(
"router"
)).
decode
()
router_mode
=
(
await
self
.
etcd_kv_cache
.
get
(
"router"
)).
decode
()
if
router_mode
==
"kv"
:
if
router_mode
==
"kv"
:
async
for
route_response
in
self
.
router
.
generate
(
router_generator
=
await
self
.
router
_client
.
generate
(
Tokens
(
tokens
=
engine_prompt
[
"prompt_token_ids"
]).
model_dump_json
()
Tokens
(
tokens
=
engine_prompt
[
"prompt_token_ids"
]).
model_dump_json
()
):
)
worker_id
,
prefix_hit_rate
=
route_response
.
split
(
"_"
)
decision
=
await
router_generator
.
__anext__
()
prefix_hit_rate
=
float
(
prefix_hit_rate
)
decision
=
decision
.
data
()
logger
.
info
(
worker_id
,
prefix_hit_rate
=
decision
.
split
(
"_"
)
f
"Worker ID:
{
worker_id
}
with estimated prefix hit rate:
{
prefix_hit_rate
}
"
prefix_hit_rate
=
float
(
prefix_hit_rate
)
)
logger
.
info
(
break
f
"Worker ID:
{
worker_id
}
with estimated prefix hit rate:
{
prefix_hit_rate
}
"
)
if
worker_id
==
""
:
if
worker_id
==
""
:
engine_generator
=
await
self
.
worker_client
.
generate
(
engine_generator
=
await
self
.
worker_client
.
generate
(
...
...
examples/tensorrt_llm/components/processor.py
View file @
c30c6990
...
@@ -52,6 +52,7 @@ class Processor(ChatProcessorMixin):
...
@@ -52,6 +52,7 @@ class Processor(ChatProcessorMixin):
self
.
remote_prefill
=
args
.
remote_prefill
self
.
remote_prefill
=
args
.
remote_prefill
self
.
router_mode
=
args
.
router
self
.
router_mode
=
args
.
router
self
.
min_workers
=
1
self
.
min_workers
=
1
self
.
args
=
args
super
().
__init__
(
engine_config
)
super
().
__init__
(
engine_config
)
...
@@ -65,6 +66,16 @@ class Processor(ChatProcessorMixin):
...
@@ -65,6 +66,16 @@ class Processor(ChatProcessorMixin):
.
endpoint
(
"generate"
)
.
endpoint
(
"generate"
)
.
client
()
.
client
()
)
)
if
self
.
args
.
router
==
"kv"
:
router_ns
,
router_name
=
Router
.
dynamo_address
()
# type: ignore
self
.
router_client
=
(
await
runtime
.
namespace
(
router_ns
)
.
component
(
router_name
)
.
endpoint
(
"generate"
)
.
client
()
)
while
len
(
self
.
worker_client
.
endpoint_ids
())
<
self
.
min_workers
:
while
len
(
self
.
worker_client
.
endpoint_ids
())
<
self
.
min_workers
:
logger
.
info
(
logger
.
info
(
f
"Waiting for workers to be ready.
\n
"
f
"Waiting for workers to be ready.
\n
"
...
@@ -88,15 +99,16 @@ class Processor(ChatProcessorMixin):
...
@@ -88,15 +99,16 @@ class Processor(ChatProcessorMixin):
worker_id
=
""
worker_id
=
""
if
self
.
router_mode
==
"kv"
:
if
self
.
router_mode
==
"kv"
:
async
for
route_response
in
self
.
router
.
generate
(
router_generator
=
await
self
.
router
_client
.
generate
(
preprocessed_request
.
tokens
.
model_dump_json
()
preprocessed_request
.
tokens
.
model_dump_json
()
):
)
worker_id
,
prefix_hit_rate
=
route_response
.
split
(
"_"
)
decision
=
await
router_generator
.
__anext__
()
prefix_hit_rate
=
float
(
prefix_hit_rate
)
decision
=
decision
.
data
()
logger
.
info
(
worker_id
,
prefix_hit_rate
=
decision
.
split
(
"_"
)
f
"Worker ID:
{
worker_id
}
with estimated prefix hit rate:
{
prefix_hit_rate
}
"
prefix_hit_rate
=
float
(
prefix_hit_rate
)
)
logger
.
info
(
break
f
"Worker ID:
{
worker_id
}
with estimated prefix hit rate:
{
prefix_hit_rate
}
"
)
if
worker_id
==
""
:
if
worker_id
==
""
:
if
self
.
router_mode
==
"round-robin"
:
if
self
.
router_mode
==
"round-robin"
:
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
.
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment