Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

perf(weave): push heavy conditions into WHERE for calls stream query #3501

Closed
wants to merge 28 commits into from
Closed
Show file tree
Hide file tree
Changes from 19 commits
Commits
Show all changes
28 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
56 changes: 56 additions & 0 deletions tests/trace/test_client_trace.py
Original file line number Diff line number Diff line change
Expand Up @@ -3300,3 +3300,59 @@ def test():
select_query["query"].count("any(calls_merged.output_dump) AS output_dump")
== 1
)


def test_calls_stream_heavy_condition_aggregation_parts(client):
def _make_query(field: str, value: str) -> tsi.CallsQueryRes:
query = {
"$in": [
{"$getField": field},
[{"$literal": value}],
]
}
res = get_client_trace_server(client).calls_query_stream(
tsi.CallsQueryReq.model_validate(
{
"project_id": get_client_project_id(client),
"query": {"$expr": query},
}
)
)
return list(res)

call_id = generate_id()
trace_id = generate_id()
parent_id = generate_id()
start = tsi.StartedCallSchemaForInsert(
project_id=client._project_id(),
id=call_id,
op_name="test_name",
trace_id=trace_id,
parent_id=parent_id,
started_at=datetime.datetime.now(tz=datetime.timezone.utc)
- datetime.timedelta(seconds=1),
attributes={"a": 5},
inputs={"param": {"value1": "hello"}},
)
client.server.call_start(tsi.CallStartReq(start=start))

res = _make_query("inputs.param.value1", "hello")
assert len(res) == 1
assert res[0].inputs["param"]["value1"] == "hello"
assert not res[0].output

end = tsi.EndedCallSchemaForInsert(
project_id=client._project_id(),
id=call_id,
ended_at=datetime.datetime.now(tz=datetime.timezone.utc),
summary={"c": 5},
output={"d": 5},
)
client.server.call_end(tsi.CallEndReq(end=end))

res = _make_query("inputs.param.value1", "hello")
assert len(res) == 1
assert res[0].inputs["param"]["value1"] == "hello"

# Does the query return the output?
assert res[0].output["d"] == 5
Copy link
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This test highlights the error case that the query creates.

33 changes: 15 additions & 18 deletions tests/trace_server/test_calls_query_builder.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,18 @@
from weave.trace_server.orm import ParamBuilder


def assert_sql(cq: CallsQuery, exp_query, exp_params):
pb = ParamBuilder("pb")
query = cq.as_sql(pb)
params = pb.get_params()

exp_formatted = sqlparse.format(exp_query, reindent=True)
found_formatted = sqlparse.format(query, reindent=True)

assert exp_formatted == found_formatted
assert exp_params == params


def test_query_baseline() -> None:
cq = CallsQuery(project_id="project")
cq.add_field("id")
Expand Down Expand Up @@ -273,15 +285,13 @@ def test_query_heavy_column_simple_filter_with_order_and_limit_and_mixed_query_c
SELECT
calls_merged.id AS id,
any(calls_merged.inputs_dump) AS inputs_dump
FROM calls_merged
FROM calls_merged FINAL
WHERE
calls_merged.project_id = {pb_2:String}
AND
(calls_merged.id IN filtered_calls)
(calls_merged.id IN filtered_calls) AND
(JSON_VALUE(calls_merged.inputs_dump, {pb_3:String}) = {pb_4:String})
GROUP BY (calls_merged.project_id, calls_merged.id)
HAVING (
JSON_VALUE(any(calls_merged.inputs_dump), {pb_3:String}) = {pb_4:String}
)
ORDER BY any(calls_merged.started_at) DESC
LIMIT 10
""",
Expand All @@ -295,19 +305,6 @@ def test_query_heavy_column_simple_filter_with_order_and_limit_and_mixed_query_c
)


def assert_sql(cq: CallsQuery, exp_query, exp_params):
pb = ParamBuilder("pb")
query = cq.as_sql(pb)
params = pb.get_params()

assert exp_params == params

exp_formatted = sqlparse.format(exp_query, reindent=True)
found_formatted = sqlparse.format(query, reindent=True)

assert exp_formatted == found_formatted


def test_query_light_column_with_costs() -> None:
cq = CallsQuery(
project_id="UHJvamVjdEludGVybmFsSWQ6Mzk1NDg2Mjc=", include_costs=True
Expand Down
147 changes: 138 additions & 9 deletions weave/trace_server/calls_query_builder.py
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,7 @@ def is_heavy(self) -> bool:

class CallsMergedAggField(CallsMergedField):
agg_fn: str
use_agg_fn: bool = True

def as_sql(
self,
Expand All @@ -81,6 +82,8 @@ def as_sql(
cast: Optional[tsi_query.CastTo] = None,
) -> str:
inner = super().as_sql(pb, table_alias)
if not self.use_agg_fn:
return clickhouse_cast(inner)
return clickhouse_cast(f"{self.agg_fn}({inner})")


Expand Down Expand Up @@ -258,11 +261,14 @@ class Condition(BaseModel):
operand: "tsi_query.Operand"
_consumed_fields: Optional[list[CallsMergedField]] = None

def as_sql(self, pb: ParamBuilder, table_alias: str) -> str:
def as_sql(
self, pb: ParamBuilder, table_alias: str, ignore_agg_fn: bool = False
) -> str:
conditions = process_query_to_conditions(
tsi_query.Query.model_validate({"$expr": {"$and": [self.operand]}}),
pb,
table_alias,
ignore_agg_fn=ignore_agg_fn,
)
if self._consumed_fields is None:
self._consumed_fields = []
Expand All @@ -283,6 +289,12 @@ def is_heavy(self) -> bool:
return True
return False

def is_feedback(self) -> bool:
for field in self._get_consumed_fields():
if isinstance(field, CallsMergedFeedbackPayloadField):
return True
return False


class HardCodedFilter(BaseModel):
filter: tsi.CallsFilter
Expand Down Expand Up @@ -448,8 +460,10 @@ def as_sql(self, pb: ParamBuilder, table_alias: str = "calls_merged") -> str:
)
SELECT {SELECT_FIELDS}
FROM calls_merged
FINAL -- optional, if heavy filter conditions
WHERE project_id = {PROJECT_ID}
AND id IN filtered_calls
AND {HEAVY_FILTER_CONDITIONS} -- optional heavy filter conditions
GROUP BY (project_id, id)
--- IF ORDER BY CANNOT BE PUSHED DOWN ---
HAVING {HEAVY_FILTER_CONDITIONS} -- optional <-- yes, this is inside the conditional
Expand Down Expand Up @@ -558,7 +572,8 @@ def as_sql(self, pb: ParamBuilder, table_alias: str = "calls_merged") -> str:
# TODO: We should unify the calls query order by fields to be orm sort by fields
order_by_fields = [
tsi.SortBy(
field=sort_by.field.field, direction=sort_by.direction.lower()
field=sort_by.field.field,
direction=cast(Literal["asc", "desc"], sort_by.direction.lower()),
)
for sort_by in self.order_fields
]
Expand Down Expand Up @@ -586,21 +601,38 @@ def _as_sql_base_format(
)

having_filter_sql = ""
having_conditions_sql: list[str] = []
having_light_conditions_sql: list[str] = []
if len(self.query_conditions) > 0:
having_conditions_sql.extend(
c.as_sql(pb, table_alias) for c in self.query_conditions
having_light_conditions_sql.extend(
c.as_sql(pb, table_alias)
for c in self.query_conditions
if not c.is_heavy() or c.is_feedback()
)
for query_condition in self.query_conditions:
for field in query_condition._get_consumed_fields():
if isinstance(field, CallsMergedFeedbackPayloadField):
needs_feedback = True
if self.hardcoded_filter is not None:
having_conditions_sql.append(self.hardcoded_filter.as_sql(pb, table_alias))
having_light_conditions_sql.append(
self.hardcoded_filter.as_sql(pb, table_alias)
)

if len(having_conditions_sql) > 0:
if len(having_light_conditions_sql) > 0:
having_filter_sql = "HAVING " + combine_conditions(
having_conditions_sql, "AND"
having_light_conditions_sql, "AND"
)

heavy_filter_sql = ""
simple_like_conditions_sql = ""
for condition in self.query_conditions:
if not condition.is_heavy() or condition.is_feedback():
continue

simple_like_conditions_sql += make_simple_like_condition(
pb, condition, table_alias
)
heavy_filter_sql += " AND " + condition.as_sql(
pb, table_alias, ignore_agg_fn=True
)

order_by_sql = ""
Expand Down Expand Up @@ -654,6 +686,8 @@ def _as_sql_base_format(
{feedback_where_sql}
{id_mask_sql}
{id_subquery_sql}
{simple_like_conditions_sql}
{heavy_filter_sql}
GROUP BY (calls_merged.project_id, calls_merged.id)
{having_filter_sql}
{order_by_sql}
Expand All @@ -664,6 +698,95 @@ def _as_sql_base_format(
return _safely_format_sql(raw_sql)


def _like_value(value: str) -> str:
"""Helper to create a LIKE pattern for string matching."""
return f"%{value}%"


def make_simple_like_condition(
pb: ParamBuilder, condition: Condition, table_alias: str = "calls_merged"
) -> str:
"""Creates a simple LIKE condition for the SQL query. Used to prepend
before JSON_VALUE is used. Optimization for ch s3 reads.

Only applies to string fields.

Args:
pb: Parameter builder for generating SQL parameters
condition: The condition to convert to SQL
table_alias: The table alias to use in the SQL (default: "calls_merged")

Returns:
SQL string with the LIKE condition or empty string if not applicable
"""
# Handle different operation types
if isinstance(condition.operand, tsi_query.EqOperation):
# Get field name from first operand
field_op = condition.operand.eq_[0]
if not isinstance(field_op, tsi_query.GetFieldOperator):
return ""
field = get_field_by_name(field_op.get_field_)
field_name = field.field

# Get value from second operand
value_op = condition.operand.eq_[1]
if not isinstance(value_op, tsi_query.LiteralOperation):
return ""
# Only apply to string values
if not isinstance(value_op.literal_, str):
return ""
field_value = value_op.literal_

elif isinstance(condition.operand, tsi_query.ContainsOperation):
# Get field name from input operand
field_op = condition.operand.contains_.input
if not isinstance(field_op, tsi_query.GetFieldOperator):
return ""
field = get_field_by_name(field_op.get_field_)
field_name = field.field

# Get value from substr operand
value_op = condition.operand.contains_.substr
if not isinstance(value_op, tsi_query.LiteralOperation):
return ""
# Only apply to string values
if not isinstance(value_op.literal_, str):
return ""
field_value = value_op.literal_

elif isinstance(condition.operand, tsi_query.InOperation):
# Get field name from first operand
field_op = condition.operand.in_[0]
if not isinstance(field_op, tsi_query.GetFieldOperator):
return ""
field = get_field_by_name(field_op.get_field_)
field_name = field.field

# Build OR conditions for each value in the IN clause
sql_parts = []
for op in condition.operand.in_[1]:
if not isinstance(op, tsi_query.LiteralOperation):
return ""
# Only apply to string values
if not isinstance(op.literal_, str):
return ""
op_wildcard_value = _like_value(op.literal_)
op_sql_value = _param_slot(pb.add_param(op_wildcard_value), "String")
sql_parts.append(f"{table_alias}.{field_name} LIKE {op_sql_value}")

if not sql_parts:
return ""
return f" AND ({' OR '.join(sql_parts)})"

else:
return ""

# For EQ and CONTAINS operations, create a single LIKE condition
wildcard_value = _like_value(field_value)
sql_value = _param_slot(pb.add_param(wildcard_value), "String")
return f" AND {table_alias}.{field_name} LIKE {sql_value}"


ALLOWED_CALL_FIELDS = {
"project_id": CallsMergedField(field="project_id"),
"id": CallsMergedField(field="id"),
Expand Down Expand Up @@ -713,6 +836,7 @@ def process_query_to_conditions(
query: tsi.Query,
param_builder: ParamBuilder,
table_alias: str,
ignore_agg_fn: bool = False,
) -> FilterToConditions:
"""Converts a Query to a list of conditions for a clickhouse query."""
conditions = []
Expand Down Expand Up @@ -781,7 +905,12 @@ def process_operand(operand: "tsi_query.Operand") -> str:
)
elif isinstance(operand, tsi_query.GetFieldOperator):
structured_field = get_field_by_name(operand.get_field_)
field = structured_field.as_sql(param_builder, table_alias)
if isinstance(structured_field, CallsMergedAggField) and ignore_agg_fn:
structured_field.use_agg_fn = False
field = structured_field.as_sql(param_builder, table_alias)
structured_field.use_agg_fn = True # reset to default
else:
field = structured_field.as_sql(param_builder, table_alias)
raw_fields_used[structured_field.field] = structured_field
return field
elif isinstance(operand, tsi_query.ConvertOperation):
Expand Down
4 changes: 2 additions & 2 deletions weave/trace_server/clickhouse_trace_server_batched.py
Original file line number Diff line number Diff line change
Expand Up @@ -1944,8 +1944,8 @@ def _ch_call_dict_to_call_schema_dict(ch_call_dict: dict) -> dict:
"op_name": ch_call_dict.get("op_name"),
"started_at": started_at,
"ended_at": ended_at,
"attributes": _dict_dump_to_dict(ch_call_dict.get("attributes_dump", "{}")),
"inputs": _dict_dump_to_dict(ch_call_dict.get("inputs_dump", "{}")),
"attributes": _dict_dump_to_dict(ch_call_dict.get("attributes_dump") or "{}"),
"inputs": _dict_dump_to_dict(ch_call_dict.get("inputs_dump") or "{}"),
"output": _nullable_any_dump_to_any(ch_call_dict.get("output_dump")),
"summary": make_derived_summary_fields(
summary=summary or {},
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
ALTER TABLE calls_merged
MODIFY SETTING min_bytes_for_wide_part = 10485760, min_age_to_force_merge_seconds = 0;
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
ALTER TABLE calls_merged
MODIFY SETTING min_bytes_for_wide_part = 0, min_age_to_force_merge_seconds = 30;
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
ALTER TABLE calls_merged
ADD INDEX bf_inputs_dump inputs_dump TYPE BLOOM_FILTER GRANULARITY 1;

ALTER TABLE calls_merged
ADD INDEX bf_output_dump output_dump TYPE BLOOM_FILTER GRANULARITY 1;

ALTER TABLE calls_merged MATERLIAZE INDEX bf_inputs_dump;
ALTER TABLE calls_merged MATERLIAZE INDEX bf_output_dump;
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
ALTER TABLE calls_merged
MODIFY SETTING min_bytes_for_wide_part = 0, min_age_to_force_merge_seconds = 30;
Loading