הזרמה דו-כיוונית ב-Agent Runtime מאפשרת תקשורת מתמשכת ודו-כיוונית בין האפליקציה לבין סוכן, ויוצאת מדפוסי בקשה-תגובה רגילים. במסמך הזה מוסבר איך לפתח, לבדוק ולפרוס סוכני סטרימינג דו-כיווני לתרחישי שימוש בזמן אמת, כמו אינטראקציה עם אודיו או וידאו.
סקירה כללית
סטרימינג דו-כיווני מספק ערוץ תקשורת מתמשך ודו-כיווני בין האפליקציה שלכם לבין הסוכן, ומאפשר לכם לחרוג מדפוסי בקשה-תגובה שמבוססים על תורות. סטרימינג דו-כיווני מתאים לתרחישי שימוש שבהם הסוכן צריך לעבד מידע ולהגיב ברציפות, כמו אינטראקציה עם קלט אודיו או וידאו עם זמן אחזור נמוך.
הזרמה דו-כיוונית עם Agent Runtime תומכת בתרחישי שימוש אינטראקטיביים של סוכנים בזמן אמת ובחילופי נתונים עבור ממשקי API של Multimodal Live. כל המסגרות תומכות בהעברת נתונים דו-כיוונית, ואפשר להשתמש בשיטות מותאמות אישית להעברת נתונים דו-כיוונית באמצעות רישום שיטות מותאמות אישית. אתם יכולים להשתמש בשידור דו-כיווני כדי ליצור אינטראקציה ישירה עם Gemini Live API או באמצעות ערכת הכלים לפיתוח סוכנים (ADK) ב-Agent Platform.
Google GenAI SDK תומך היטב בפריסה של סוכן מרוחק עם שיטות שאילתה דו-כיווניות. כדי לפרוס סוכן עם יכולת דו-כיוונית, צריך להגדיר את EXPERIMENTAL agent server mode כשמשתמשים ב-SDK או כשמבצעים קריאה ל-Agent Platform API.
פיתוח סוכן
כדי להטמיע סטרימינג דו-כיווני בזמן פיתוח סוכן:
רישום של שיטות מותאמות אישית (אופציונלי)
הגדרת שיטת שאילתה להעברה דו-כיוונית בסטרימינג
כדי להפוך את הסוכן ל "תומך בדו-כיווניות", צריך להגדיר method bidi_stream_query שמקבלת בקשות לזרם כקלט ומחזירה תגובות לזרם כפלט באופן אסינכרוני. לדוגמה, התבנית הבאה מרחיבה את התבנית הבסיסית כדי להזרים בקשות ותשובות, ואפשר לפרוס אותה ב-Gemini Enterprise Agent Platform:
import asyncio
from typing import Any, AsyncIterable
class BidiStreamingAgent(StreamingAgent):
async def bidi_stream_query(
self,
request_queue: asyncio.Queue[Any]
) -> AsyncIterable[Any]:
from langchain.load.dump import dumpd
while True:
request = await request_queue.get()
# This is just an illustration, you're free to use any termination mechanism.
if request == "END":
break
for chunk in self.graph.stream(request):
yield dumpd(chunk)
agent = BidiStreamingAgent(
model=model, # Required.
tools=[get_exchange_rate], # Optional.
project="PROJECT_ID",
location="LOCATION",
)
agent.set_up()
כשמשתמשים ב-API לסטרימינג דו-כיווני, חשוב לזכור את הנקודות הבאות:
asyncio.Queue: אפשר להוסיף לתור הבקשות הזה נתונים מכל סוג, כדי להמתין לשליחה אל ה-Model API.זמן קצוב מקסימלי: הזמן הקצוב המקסימלי לשאילתת סטרימינג דו-כיווני הוא 10 דקות. אם הסוכן שלכם דורש זמני עיבוד ארוכים יותר, כדאי לחלק את המשימה לחלקים קטנים יותר ולהשתמש בסשן או בזיכרון כדי לשמור את המצב.
הגבלת צריכת התוכן: כשצורכים תוכן ממקור נתונים דו-כיווני, חשוב לנהל את הקצב שבו הסוכן מעבד את הנתונים הנכנסים. אם הסוכן צורך נתונים לאט מדי, זה עלול לגרום לבעיות כמו זמן אחזור מוגבר או עומס על הזיכרון בצד השרת. מטמיעים מנגנונים לשליפת נתונים באופן פעיל כשהסוכן מוכן לעבד אותם, ונמנעים מחסימת פעולות שעלולות להפסיק את צריכת התוכן.
הגבלת קצב יצירת התוכן: אם נתקלים בבעיות של לחץ חוזר (back pressure) (כשהמפיק יוצר נתונים מהר יותר מהקצב שבו הצרכן יכול לעבד אותם), צריך להגביל את קצב יצירת התוכן. הפעולה הזו יכולה לעזור למנוע הצפת חוצץ ולהבטיח חוויית סטרימינג חלקה.
בדיקת השיטה של שאילתות להעברה דו-כיוונית בסטרימינג
כדי לבדוק את השאילתה של הזרמת נתונים דו-כיוונית באופן מקומי, קוראים לשיטה bidi_stream_query ומבצעים איטרציה על התוצאות:
import asyncio
import pprint
import time
request_queue = asyncio.Queue()
async def generate_input():
# This is just an illustration, you're free to use any appropriate input generator.
request_queue.put_nowait(
{"input": "What is the exchange rate from US dolloars to Swedish currency"}
)
time.sleep(5)
request_queue.put_nowait(
{"input": "What is the exchange rate from US dolloars to Euro currency"}
)
time.sleep(5)
request_queue.put_nowait("END")
async def print_query_result():
async for chunk in agent.bidi_stream_query(request_queue):
pprint.pprint(chunk, depth=1)
input_task = asyncio.create_task(generate_input())
output_task = asyncio.create_task(print_query_result())
await asyncio.gather(input_task, output_task, return_exceptions=True)
אותו חיבור שאילתה דו-כיווני יכול לטפל בכמה בקשות ותשובות. לכל בקשה חדשה בתור, הדוגמה הבאה יוצרת זרם של נתחים שמכילים מידע שונה על התגובה:
{'actions': [...], 'messages': [...]}
{'messages': [...], 'steps': [...]}
{'messages': [...], 'output': 'The exchange rate from US dollars to Swedish currency is 1 USD to 10.5751 SEK. \n'}
{'actions': [...], 'messages': [...]}
{'messages': [...], 'steps': [...]}
{'messages': [...], 'output': 'The exchange rate from US dollars to Euro currency is 1 USD to 0.86 EUR. \n'}
אופציונלי: רישום של שיטות מותאמות אישית
אפשר לרשום פעולות במצבי ביצוע רגילים (מיוצגים על ידי מחרוזת ריקה ""), סטרימינג (stream) או סטרימינג דו-כיווני (bidi_stream).
from typing import AsyncIterable, Iterable
class CustomAgent(BidiStreamingAgent):
# ... same get_state and get_state_history function definition.
async def get_state_bidi_mode(
self,
request_queue: asyncio.Queue[Any]
) -> AsyncIterable[Any]:
while True:
request = await request_queue.get()
if request == "END":
break
yield self.graph.get_state(request)._asdict()
def register_operations(self):
return {
# The list of synchrounous operations to be registered
"": ["query", "get_state"]
# The list of streaming operations to be registered
"stream": ["stream_query", "get_state_history"]
# The list of bidi streaming operations to be registered
"bidi_stream": ["bidi_stream_query", "get_state_bidi_mode"]
}
פריסת סוכן
אחרי שמפתחים את הסוכן כ-live_agent, אפשר לפרוס את הסוכן ב-Agent Platform על ידי יצירת מופע של Agent Platform.
שימו לב: באמצעות GenAI SDK, כל הגדרות הפריסה (חבילות נוספות ואמצעי בקרה מותאמים אישית על משאבים) מוקצות כערך של config כשיוצרים את מופע Agent Platform.
מאתחלים את לקוח ה-GenAI:
import vertexai
from vertexai import types as vertexai_types
client = vertexai.Client(project=PROJECT, location=LOCATION)
פורסים את הסוכן ב-Agent Platform. שימו לב שהתג EXPERIMENTAL
agent_server_mode נדרש לסוכן שתומך בסטרימינג דו-כיווני:
remote_live_agent = client.agent_engines.create(
agent=live_agent,
config={
"staging_bucket": STAGING_BUCKET,
"requirements": [
"google-cloud-aiplatform[agent_engines,adk]==1.88.0",
"cloudpickle==3.0",
"websockets"
],
"agent_server_mode": vertexai_types.AgentServerMode.EXPERIMENTAL,
},
)
במאמר יצירת מופע של Agent Runtime מוסבר על השלבים שמתבצעים ברקע במהלך הפריסה.
מקבלים את מזהה המשאב של הסוכן:
remote_live_agent.api_resource.name
שימוש בסוכן
אם הגדרתם פעולה מסוג bidi_stream_query כשפיתחתם את הסוכן, אתם יכולים להשתמש ב-GenAI SDK ל-Python כדי להריץ שאילתות אסינכרוניות לסוכן באמצעות סטרימינג דו-כיווני.
אפשר לשנות את הדוגמה הבאה עם כל נתון שהסוכן יכול לזהות, באמצעות כל לוגיקת סיום רלוונטית לזרם הקלט ולזרם הפלט:
async with client.aio.live.agent_engines.connect(
agent_engine=remote_live_agent.api_resource.name,
config={"class_method": "bidi_stream_query"}
) as connection:
while True:
#
input_str = input("Enter your question: ")
if input_str == "exit":
break
await connection.send({"input": input_str})
while True:
response = await connection.receive()
print(response)
if response["bidiStreamOutput"]["output"] == "end of turn":
break
Agent Runtime מעביר תשובות כרצף של אובייקטים שנוצרים באופן איטרטיבי. לדוגמה, קבוצה של שתי תגובות בתור הראשון עשויה להיראות כך:
Enter your next question: Weather in San Diego?
{'bidiStreamOutput': {'output': "FunctionCall: {'name': 'get_current_weather', 'args': {'location': 'San Diego'}}\n"}}
{'bidiStreamOutput': {'output': 'end of turn'}}
Enter your next question: exit
שימוש בסוכן של ערכת פיתוח סוכנים
אם פיתחתם את הסוכן באמצעות ערכה לפיתוח סוכנים (ADK), אתם יכולים להשתמש בסטרימינג דו-כיווני כדי ליצור אינטראקציה עם Gemini Live API.
בדוגמה הבאה נוצר סוכן שיחה שמקבל שאלות טקסט של משתמשים ונתוני אודיו של תשובות מ-Gemini Live API:
import numpy as np
from google.adk.agents.live_request_queue import LiveRequest
from google.adk.events import Event
from google.genai import types
def prepare_live_request(input_text: str) -> LiveRequest:
part = types.Part.from_text(text=input_text)
content = types.Content(parts=[part])
return LiveRequest(content=content)
async with client.aio.live.agent_engines.connect(
agent_engine=remote_live_agent.api_resource.name,
config={
"class_method": "bidi_stream_query",
"input": {"input_str": "hello"},
}) as connection:
first_req = True
while True:
input_text = input("Enter your question: ")
if input_text == "exit":
break
if first_req:
await connection.send({
"user_id": USER_ID,
"live_request": prepare_live_request(input_text).dict()
})
first_req = False
else:
await connection.send(prepare_live_request(input_text).dict())
audio_data = []
while True:
async def receive():
return await connection.receive()
receiving = asyncio.Task(receive())
done, _ = await asyncio.wait([receiving])
if receiving not in done:
receiving.cancel()
break
event = Event.model_validate(receiving.result()["bidiStreamOutput"])
part = event.content and event.content.parts and event.content.parts[0]
if part.inline_data and part.inline_data.data:
chunk_data = part.inline_data.data
data = np.frombuffer(chunk_data, dtype=np.int16)
audio_data.append(data)
else:
print(part)
if audio_data:
concatenated_audio = np.concatenate(audio_data)
display(Audio(concatenated_audio, rate=24000, autoplay=True))