CrewAI lets you build multi-agent systems where specialized agents collaborate on tasks. This guide shows how to give CrewAI agents access to Catalogian tools for data source monitoring and analysis.
cat_live_): available once any record is on a paid tier (Watch or Max)pip install crewai openaiFirst, create tool functions that call Catalogian's Responses API:
import os
import json
from openai import OpenAI
from crewai.tools import tool
catalogian = OpenAI(
base_url="https://api.catalogian.com/v1",
api_key=os.environ["CATALOGIAN_KEY"],
)
def _call_catalogian(tool_name: str, args: dict = {}) -> str:
response = catalogian.responses.create(
model="catalogian-1",
input=json.dumps(args) if args else "{}",
tools=[{"type": "function", "name": tool_name}],
tool_choice={"type": "function", "name": tool_name},
)
for item in response.output:
if hasattr(item, "content"):
for part in item.content:
if hasattr(part, "text"):
return part.text
return "No result"
@tool("list_records")
def list_records() -> str:
"""List all records in Catalogian."""
return _call_catalogian("list_records")
@tool("get_delta")
def get_delta(record_slug: str, limit: int = 5) -> str:
"""Get recent change events for a record."""
return _call_catalogian("get_delta", {
"recordSlug": record_slug, "limit": limit
})
@tool("profile_snapshot")
def profile_snapshot(record_slug: str) -> str:
"""Get field stats, null rates, and data quality metrics for a record."""
return _call_catalogian("profile_snapshot", {
"recordSlug": record_slug
})
@tool("search_snapshot")
def search_snapshot(record_slug: str, query: str) -> str:
"""Search across all fields in a record."""
return _call_catalogian("search_snapshot", {
"recordSlug": record_slug, "query": query
})
@tool("get_delta_rows")
def get_delta_rows(record_slug: str, delta_event_id: str, change_type: str = "changed") -> str:
"""Get row-level before/after data for a specific delta event."""
return _call_catalogian("get_delta_rows", {
"recordSlug": record_slug,
"deltaEventId": delta_event_id,
"changeType": change_type,
})Create specialized agents for different data monitoring tasks:
from crewai import Agent
catalog_monitor = Agent(
role="Data Source Monitor",
goal="Monitor data sources for significant changes and anomalies",
backstory="""You are a data specialist who monitors important sources
for changes. You detect new rows, changed values, deleted rows,
and data quality issues. Always call profile_snapshot before analyzing
a new record to understand its structure.""",
tools=[list_records, get_delta, profile_snapshot, get_delta_rows],
verbose=True,
)
data_analyst = Agent(
role="Data Analyst",
goal="Analyze record data and provide actionable insights",
backstory="""You analyze record data to find trends, anomalies,
and opportunities. You search for specific rows and compare data
across time periods.""",
tools=[search_snapshot, profile_snapshot, get_delta],
verbose=True,
)from crewai import Crew, Task
# Define tasks
monitor_task = Task(
description="""Check all records for recent changes.
For any record with changes in the last 24 hours:
1. Get the delta events
2. Summarize what changed (new rows, changed values, deletions)
3. Flag anything unusual (large deletions, many changed rows)""",
agent=catalog_monitor,
expected_output="A summary of recent changes across all records",
)
analysis_task = Task(
description="""Based on the monitoring report, analyze the most significant
changes in detail. For records with changed rows, compare before/after
values. Identify any patterns or concerns.""",
agent=data_analyst,
expected_output="Detailed analysis with specific examples and recommendations",
)
# Assemble and run
crew = Crew(
agents=[catalog_monitor, data_analyst],
tasks=[monitor_task, analysis_task],
verbose=True,
)
result = crew.kickoff()
print(result)For simpler use cases, a single agent works well:
from crewai import Agent, Task, Crew
agent = Agent(
role="Data Assistant",
goal="Answer questions about record data",
backstory="You help users query and understand their data sources.",
tools=[list_records, get_delta, profile_snapshot, search_snapshot],
)
task = Task(
description="Find rows mentioning 'Cobb' in west-coast-earthquakes "
"and check if any changed this week.",
agent=agent,
expected_output="List of matching rows with any recent changes",
)
crew = Crew(agents=[agent], tasks=[task])
result = crew.kickoff()
print(result)Rate limits: Each tool call makes one request to Catalogian. The MCP/Responses endpoint allows 30 requests per minute per key. For multi-agent crews with many tool calls, consider adding delays between tasks or using separate API keys per agent.
Tips for building effective AI agent integrations. Agent Best Practices →