AI Agent Production Deployment and Engineering
Production deployment involves multiple aspects such as reliability, scalability, and cost control.
Engineering practices ensure stable system operation and ease of maintenance.
Production Deployment Architecture
Production environments place higher demands on the system.
Multiple aspects such as reliability, scalability, monitoring, and alerting need to be considered.
Deployment Modes
Single-machine Deployment
Suitable for development and testing environments.
Simple and easy to deploy, but cannot handle production-level traffic.
Distributed Deployment
Suitable for production environments; multiple aspects need to be considered.
Load Balancing: Distribute requests to multiple instances.
Service Discovery: Dynamically manage the list of instances.
State Management: Handle distributed state consistency issues.
Fault Tolerance: A single point of failure does not affect the overall service.
Containerized Deployment
Dockerfile Example
FROM python:3.11-slim
# Set working directory
WORKDIR /app
# Install dependencies
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# Copy application code
COPY . .
# Create non-root user (security consideration)
RUN useradd -m appuser && chown -R appuser:appuser /app
USER appuser
# Expose port
EXPOSE 8000
# Startup command
CMD ["python", "agent.py"]
Kubernetes Deployment
Kubernetes Deployment Configuration
apiVersion: apps/v1
kind: Deployment
metadata:
name: ai-agent
labels:
app: ai-agent
spec:
# Replica Count
replicas: 3
selector:
matchLabels:
app: ai-agent
template:
metadata:
labels:
app: ai-agent
spec:
containers:
- name: agent
image: your-registry/ai-agent:latest
ports:
- containerPort: 8000
# Resource Limits
resources:
limits:
memory: "2Gi"
cpu: "1"
requests:
memory: "1Gi"
cpu: "0.5"
# Health Check
livenessProbe:
httpGet:
path: /health
port: 8000
initialDelaySeconds: 10
periodSeconds: 30
readinessProbe:
httpGet:
path: /ready
port: 8000
initialDelaySeconds: 5
periodSeconds: 10
---
apiVersion: v1
kind: Service
metadata:
name: ai-agent-service
spec:
selector:
app: ai-agent
ports:
- port: 80
targetPort: 8000
type: LoadBalancer
Cost Control
The cost of an Agent system mainly comes from LLM API calls.
Optimizing cost is an important topic for production deployment.
Token Optimization Strategies
Prompt Compression: Simplify the prompt and reduce unnecessary text.
Context Truncation: Keep only the key context content.
Response Caching: Cache responses for identical or similar queries.
Code Implementation
Cost Optimization Agent Implementation
import json
from functools import lru_cache
class CostOptimizedAgent:
"""
Cost-Optimized Agent
Reduce API call costs through caching and compression
"""
def __init__(self, base_agent, cache, embedder=None):
# Base Agent
self.base_agent = base_agent
# Cache Storage
self.cache = cache
# Embedder (for similarity cache matching)
self.embedder = embedder
def process(self, request):
"""
Handle requests, including cost optimization logic
"""
# Generate request fingerprint
request_fingerprint = self.generate_fingerprint(request)
# Check exact cache
cached_result = self.cache.get(request_fingerprint)
if cached_result:
return {
**cached_result,
"from_cache": True
}
# If there is an embedder, check semantic similarity cache
if self.embedder:
similar = self.find_similar_cached(request)
if similar:
return {
**similar,
"from_cache": True,
"similar_to": similar.get("request_id")
}
# Execute request
result = self.base_agent.process(request)
# Cache result
self.cache.set(
request_fingerprint,
result,
ttl=3600 # Expire in 1 hour
)
return {
**result,
"from_cache": False
}
def generate_fingerprint(self, request):
"""
Generate request fingerprint
For exact cache matching
"""
content = json.dumps(request, sort_keys=True)
return hashlib.sha256(content.encode()).hexdigest()
def find_similar_cached(self, request, threshold=0.95):
"""
Find semantically similar cached results
Using cosine similarity of embedding vectors
"""
if not self.embedder:
return None
# Vectorize the request
request_vector = self.embedder.embed([str(request)])[0]
# Search for similar requests in the cache
best_match = None
best_similarity = 0
for cached in self.cache.get_all():
cached_vector = cached["vector"]
similarity = cosine_similarity(request_vector, cached_vector)
if similarity > best_similarity and similarity >= threshold:
best_similarity = similarity
best_match = cached
return best_match
class ResponseCache:
"""
Response Cache
Simple in-memory cache implementation
Redis can be used in production
"""
def __init__(self):
self.cache = {}
self.timestamps = {}
def get(self, key):
"""Get cache"""
if key in self.cache:
# Check if expired
if self.is_expired(key):
self.delete(key)
return None
return self.cache[key]
return None
def set(self, key, value, ttl=3600):
"""Set cache"""
self.cache[key] = value
self.timestamps[key] = time.time() + ttl
def delete(self, key):
"""Delete cache"""
if key in self.cache:
del self.cache[key]
if key in self.timestamps:
del self.timestamps[key]
def is_expired(self, key):
"""Check if expired"""
if key not in self.timestamps:
return True
return time.time() > self.timestamps[key]
def get_all(self):
"""Get all cache items"""
return list(self.cache.values())
def cosine_similarity(vec1, vec2):
"""Calculate cosine similarity"""
dot_product = sum(a * b for a, b in zip(vec1, vec2))
norm1 = math.sqrt(sum(a * a for a in vec1))
norm2 = math.sqrt(sum(b * b for b in vec2))
return dot_product / (norm1 * norm2)
Streaming and Asynchronous
Streaming Response
Streaming responses reduce wait time and improve user experience.
You don't have to wait for the full response; you can display generated content incrementally.
Code Implementation
Streaming Response Agent
class StreamingAgent:
"""
Agent that supports streaming responses
Return while generating, reducing wait time
"""
def __init__(self, llm):
self.llm = llm
async def stream_generate(self, prompt):
"""
Streaming response generation
Use async generator pattern
:yield: The generated text fragment
"""
# Start asynchronous generation task
async for chunk in self.llm.stream_generate(prompt):
yield chunk
async def process_stream(self, request):
"""
Handle streaming requests
Return async generator
"""
prompt = self.build_prompt(request)
# Return generator
async def generate():
async for chunk in self.stream_generate(prompt):
yield chunk
return generate()
# Usage example
async def main():
agent = StreamingAgent(llm)
# Start streaming generation
async for token in agent.stream_generate("Explain quantum computing"):
# Print while generating
print(token, end="", flush=True)
print() # newline
# run
asyncio.run(main())
Asynchronous Task Queue
For long-running tasks, use asynchronous queue processing.
Return immediately after the user submits a task, and execute asynchronously in the background.
Asynchronous Task Queue Implementation
from queue import Queue
from threading import Thread
import uuid
class AsyncAgent:
"""
Async Agent
Use a task queue to handle long-running tasks
"""
def __init__(self, agent, task_queue):
# Base Agent
self.agent = agent
# Task Queue
self.task_queue = task_queue
# Result Storage
self.results = {}
# Start background worker thread
self.worker_thread = Thread(target=self.process_queue, daemon=True)
self.worker_thread.start()
async def submit(self, task):
"""
submit task
Return task ID immediately
"""
# Generate a unique task ID
task_id = str(uuid.uuid4())
# Add to queue
await self.task_queue.enqueue({
"id": task_id,
"task": task,
"status": "pending",
"created_at": time.time()
})
return task_id
async def get_result(self, task_id):
"""
Get task results
non-blocking
"""
if task_id in self.results:
return self.results[task_id]
# Check task status
status = await self.task_queue.get_status(task_id)
if status is None:
return None # Task does not exist
if status == "pending" or status == "processing":
return {
"status": status,
"result": None
}
return None
def process_queue(self):
"""
Background worker thread
Take a task from the queue and execute it
"""
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
while True:
task = loop.run_until_complete(self.task_queue.dequeue())
if task:
# Update status to processing
loop.run_until_complete(
self.task_queue.update_status(task["id"], "processing")
)
try:
# Execute Task
result = self.agent.process(task["task"])
# Save Result
self.results[task["id"]] = {
"status": "completed",
"result": result,
"completed_at": time.time()
}
loop.run_until_complete(
self.task_queue.update_status(task["id"], "completed")
)
except Exception as e:
self.results[task["id"]] = {
"status": "failed",
"error": str(e)
}
loop.run_until_complete(
self.task_queue.update_status(task["id"], "failed")
)
class TaskQueue:
"""Simple task queue implementation"""
def __init__(self):
self.queue = Queue()
self.statuses = {}
async def enqueue(self, task):
"""Add task to queue"""
self.queue.put(task)
self.statuses[task["id"]] = "pending"
async def dequeue(self):
Retrieve a task
if not self.queue.empty():
return self.queue.get()
return None
async def get_status(self, task_id):
Get task status
return self.statuses.get(task_id)
async def update_status(self, task_id, status):
Update task status
self.statuses[task_id] = status
Microservices Design
Split the Agent system into multiple independent services to improve maintainability and scalability.
Service Partitioning
| Services | Responsibilities |
|---|---|
| API Gateway | Unified entry point, authentication rate limiting, route forwarding. |
| Agent Service | Core logic, reasoning and planning. |
| tool service | External API integration, tool invocation. |
| storage service | Knowledge base, state management, caching. |
| monitoring service | Log collection, metric monitoring, alerting. |
Inter-service Communication
synchronous communication: HTTP/gRPC, used for real-time request-response.
asynchronous communication: Message queue, used for time-consuming operations and event notifications.
Chapter Summary
This chapter introduces key practices for production deployment and engineering.
Production deployment architectureFrom standalone to distributed, scale gradually.
Containerized deploymentUse Docker and Kubernetes.
Cost controlReduce costs through caching and token optimization.
Streaming and AsynchronousImprove user experience and system throughput.
Microservice-oriented designImprove the maintainability and scalability of the system.
Production deployment requires comprehensive consideration of multiple aspects such as reliability, performance, and cost.
It is recommended to start with a small scale and gradually increase the complexity.
Other extensions