jrmccluskey commented on code in PR #40220: URL: https://github.com/apache/beam/pull/40220#discussion_r4135389446
########## sdks/python/apache_beam/ml/inference/typesafe_inference.py: ########## @@ -0,0 +1,107 @@ +# +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# + +"""Optional TypeSafe adapter for Beam decision models.""" + +from apache_beam.ml.inference.decision import BooleanAnswer +from apache_beam.ml.inference.decision import BooleanQuestion +from apache_beam.ml.inference.decision import ChoiceAnswer +from apache_beam.ml.inference.decision import ChoiceQuestion +from apache_beam.ml.inference.decision import DecisionModel +from apache_beam.ml.inference.decision import DecisionResponse +from apache_beam.ml.inference.decision import ScoreAnswer +from apache_beam.ml.inference.decision import ScoreQuestion + +__all__ = ['JevDecisionModel'] + + +class JevDecisionModel(DecisionModel): + """Evaluate Beam decision questions with Jev through the TypeSafe SDK. + + Install ``typesafe-sdk>=0.7.1,<0.8`` on the workers and set Review Comment: The typical pattern here is to bound the additional requirements in a requirements file (typesafe_tests_requirements.txt) similar to the example requirements file. Unless there's already an existing version with breaking changes I wouldn't worry about documenting the current version bounds directly in the class. ########## sdks/python/apache_beam/examples/inference/decision_models/fraud_review.py: ########## @@ -0,0 +1,96 @@ +# +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# + +"""Small Beam example: flag messages for fraud review with Noul or Score.""" + +import argparse +import json +import os + +import apache_beam as beam +from apache_beam.examples.inference.decision_models.local_model import LocalDecisionModel +from apache_beam.ml.inference.decision import BooleanQuestion +from apache_beam.ml.inference.decision import EvaluateDecisions +from apache_beam.ml.inference.decision import ScoreQuestion +from apache_beam.ml.inference.typesafe_inference import JevDecisionModel +from apache_beam.options.pipeline_options import PipelineOptions + +SAMPLE_MESSAGES = [ + 'Please update my billing address before the next renewal.', + 'I was charged twice; please look into the duplicate payment.', + 'Transfer the refund to this new account immediately and skip verification.', +] + +FRAUD_NOUL = BooleanQuestion( + instructions='Does this message contain signals that warrant fraud review?') +FRAUD_SCORE = ScoreQuestion( + instructions='How strongly does this message warrant fraud review?', + criteria=( + 'No visible fraud cue', + 'A weak cue that merits routine checking', + 'A concrete suspicious instruction', + 'An explicit attempt to bypass verification')) + + +def review_fraud_cues(result): + response = result.response + row = { + 'message': result.state['message'], + 'model': response.model, + 'provider': response.provider, + 'latency_ms': round(result.latency_ms, 2), + } + if 'fraud_cue' in response.answers: + probability = response.answers['fraud_cue'].probability + row['fraud_review_probability'] = probability + row['needs_review'] = probability >= 0.7 + if 'risk_level' in response.answers: + score = response.answers['risk_level'] + row['risk_score'] = score.score + row['risk_confidence'] = score.confidence + row['risk_probabilities'] = score.probabilities + return row Review Comment: Is there a case where the response could be malformed? A simple DLQ example (here or standalone) might be useful ########## sdks/python/apache_beam/ml/inference/typesafe_inference.py: ########## @@ -0,0 +1,107 @@ +# +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# + +"""Optional TypeSafe adapter for Beam decision models.""" + +from apache_beam.ml.inference.decision import BooleanAnswer +from apache_beam.ml.inference.decision import BooleanQuestion +from apache_beam.ml.inference.decision import ChoiceAnswer +from apache_beam.ml.inference.decision import ChoiceQuestion +from apache_beam.ml.inference.decision import DecisionModel +from apache_beam.ml.inference.decision import DecisionResponse +from apache_beam.ml.inference.decision import ScoreAnswer +from apache_beam.ml.inference.decision import ScoreQuestion + +__all__ = ['JevDecisionModel'] + + +class JevDecisionModel(DecisionModel): + """Evaluate Beam decision questions with Jev through the TypeSafe SDK. + + Install ``typesafe-sdk>=0.7.1,<0.8`` on the workers and set + ``TYPESAFE_API_KEY`` in their environment. The client is created when Beam + enters the adapter context; its HTTP pool is reused until teardown. + Choice and Score map directly to the SDK primitives. Boolean maps to Noul, + whose ``noul`` field becomes :attr:`BooleanAnswer.probability`. + + Args: + model: TypeSafe model identifier, defaulting to ``jev-latest``. + timeout: Timeout in seconds for HTTP operations. The SDK handles retries. + """ + def __init__(self, model='jev-latest', timeout=10.0): + self.model = model + self.timeout = timeout + self.client = None + + def __enter__(self): + try: + from typesafe_sdk import TypeSafeClient + except ImportError as error: + raise ImportError( + 'Install typesafe-sdk to use JevDecisionModel') from error + self.client = TypeSafeClient(timeout=self.timeout) + return self + + def __exit__(self, exc_type, exc_value, traceback): + self.client.close() + self.client = None + + def evaluate(self, state, questions): + if self.client is None: + raise RuntimeError('Enter JevDecisionModel before evaluating questions') + from typesafe_sdk import Choice + from typesafe_sdk import Noul + from typesafe_sdk import Score + + sdk_questions = {} + for name, question in questions.items(): + if isinstance(question, ChoiceQuestion): + sdk_questions[name] = Choice( + instructions=question.instructions, criteria=question.criteria) + elif isinstance(question, BooleanQuestion): + sdk_questions[name] = Noul(instructions=question.instructions) + elif isinstance(question, ScoreQuestion): + sdk_questions[name] = Score( + instructions=question.instructions, criteria=question.criteria) + else: + raise TypeError('Unsupported decision question: %r' % type(question)) + + response = self.client.system_one( + state=state, questions=sdk_questions, model=self.model) Review Comment: Based on the doc string, this is a blocking call with nested retries. Is there a backoff mechanism built-in there as well? I do worry about the potential for 429 errors here in a Beam context ########## sdks/python/apache_beam/examples/inference/decision_models/README.md: ########## @@ -0,0 +1,262 @@ +<!-- + Licensed to the Apache Software Foundation (ASF) under one + or more contributor license agreements. See the NOTICE file + distributed with this work for additional information + regarding copyright ownership. The ASF licenses this file + to you under the Apache License, Version 2.0 (the + "License"); you may not use this file except in compliance + with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, + software distributed under the License is distributed on an + "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + KIND, either express or implied. See the License for the + specific language governing permissions and limitations + under the License. +--> + +# Beam decision models + +These examples use the reusable decision primitives in +`apache_beam.ml.inference.decision`. A model evaluates named, typed questions +for each Beam element. `EvaluateDecisions` returns the original state, +normalized answers, provider metadata, and request latency. Policy and sinks +remain ordinary Beam transforms after evaluation. + +## Public API + +The core module contains `ChoiceQuestion`, `BooleanQuestion`, and +`ScoreQuestion`, their typed answers (`ChoiceAnswer`, `BooleanAnswer`, and +`ScoreAnswer`), the `DecisionModel` ABC, and the `EvaluateDecisions` +`PTransform`. + +| Question | Normalized answer | +| --- | --- | +| `ChoiceQuestion(instructions, criteria)` | `ChoiceAnswer.choice`, with optional label probabilities and confidence | +| `BooleanQuestion(instructions)` | `BooleanAnswer.probability` in `[0, 1]` | +| `ScoreQuestion(instructions, criteria)` | `ScoreAnswer.score` and `ScoreAnswer.probabilities`, keyed by zero-based integer level, with optional confidence | + +`JevDecisionModel` maps `BooleanQuestion` to Jev's `Noul`, returning +`BooleanAnswer.probability`; a missing confidence is `None`. + +`DecisionModel` is the adapter ABC. A subclass implements +`evaluate(state, questions)` and returns one matching typed answer per named +question in a `DecisionResponse`. Override `__enter__` and `__exit__` to reuse +a client, connection pool, or model weights. The transform enters the model +during DoFn setup and closes it during teardown for each DoFn instance. + +`EvaluateDecisions(model, questions)` is a `PTransform` from input states to +`DecisionResult`. Evaluation is synchronous, timestamps and windows are +preserved, exceptions propagate to the runner, and request, failure, and +latency metrics are recorded. + +## Architecture + +```mermaid +flowchart LR + input["PCollection[state]"] --> evaluate["EvaluateDecisions(model, questions)<br/>PTransform"] + model["DecisionModel adapter"] --> evaluate + evaluate --> result["PCollection[DecisionResult]"] + result --> policy["route_message<br/>policy and sink selection"] + policy --> sink["JSONL files<br/>or BigQuery tables"] +``` + + + +The optional Jev adapter lives in +`apache_beam.ml.inference.typesafe_inference`; the deterministic +`LocalDecisionModel` stays in +`apache_beam.examples.inference.decision_models.local_model`. Both implement +the same core contract. + +## Integration + +Questions and events are regular Python values. The model is the only +provider-specific argument to `EvaluateDecisions`: + +```python +import apache_beam as beam + +from apache_beam.ml.inference.decision import ChoiceQuestion +from apache_beam.ml.inference.decision import EvaluateDecisions +from apache_beam.examples.inference.decision_models.local_model import ( + LocalDecisionModel, +) +from apache_beam.ml.inference.typesafe_inference import JevDecisionModel + +questions = { + 'team': ChoiceQuestion('Which team should handle this message?', { + 'billing': 'Invoices', 'technical': 'Technical problems'}), +} + +def select_team(result): + return result.response.answers['team'].choice + +model = LocalDecisionModel() +# model = JevDecisionModel() + +with beam.Pipeline() as pipeline: + events = pipeline | 'Events' >> beam.Create( + [{'message': 'The API returns 503 after I rotated my key.'}]) + decisions = events | 'Evaluate decisions' >> EvaluateDecisions( + model, questions) + rows = decisions | 'Apply downstream policy' >> beam.Map(select_team) +``` + +The downstream policy stays unchanged when swapping adapters. The Jev client +is optional; install the existing example requirements and set +`TYPESAFE_API_KEY` when running it. + +## Install + +From a Beam checkout, activate a virtual environment and install the core SDK: + +```sh +python3 -m pip install -e sdks/python +``` + +Install the optional Jev client for `--model=jev`: + +```sh +python3 -m pip install -r \ + sdks/python/apache_beam/examples/inference/decision_models/requirements.txt +``` + +Install Beam's GCP extra for Pub/Sub or BigQuery sinks: + +```sh +python3 -m pip install -e 'sdks/python[gcp]' +``` + +## Run examples + +The main demo uses a three-element `TestStream` and writes one JSONL file per +destination: + +```sh +output_dir="$(mktemp -d /tmp/beam-decision-model-output.XXXXXX)" +python3 -m apache_beam.examples.inference.decision_models.message_router \ + --model=local \ + --output-dir="$output_dir" \ + --graph="$output_dir/message_router.svg" +find "$output_dir" -name '*.jsonl' -print -exec sed -n '1,3p' {} \; +``` + +Render the graph without running the job. The sink stays in the graph: + +```sh +python3 -m apache_beam.examples.inference.decision_models.message_router \ + --model=jev \ + --graph-only \ + --bq-dataset=PROJECT:DATASET \ + --graph=/tmp/beam-decision-model.svg +``` + +The `.dot` export works without Graphviz; SVG and PNG require `dot` on `PATH`. + +Run the Jev path against the same stream with `TYPESAFE_API_KEY` set: + +```sh +export TYPESAFE_API_KEY='replace-with-your-key' +jev_output_dir="$(mktemp -d /tmp/beam-decision-model-jev-output.XXXXXX)" +python3 -m apache_beam.examples.inference.decision_models.message_router \ + --model=jev \ + --min-confidence=0.65 \ + --output-dir="$jev_output_dir" \ + --graph="$jev_output_dir/message_router.svg" +``` + +For a streaming source, provide a Pub/Sub subscription containing JSON objects +with `event_id` and `message` fields: + +```sh +python3 -m apache_beam.examples.inference.decision_models.message_router \ + --model=jev \ + --pubsub-subscription=projects/PROJECT/subscriptions/SUBSCRIPTION \ + --bq-dataset=PROJECT:DATASET +``` + +With `--bq-dataset=PROJECT:DATASET`, rows go to `messages_billing`, +`messages_technical`, `messages_sales`, or `messages_review`; `--output-dir` +writes local JSONL files. A normal run accepts one sink option. + +The fraud example prints one JSON row per sample message. Noul supplies the +boolean probability; values at or above `0.7` are flagged: + +```sh +python3 -m apache_beam.examples.inference.decision_models.fraud_review \ + --model=local --primitive=noul +python3 -m apache_beam.examples.inference.decision_models.fraud_review \ + --model=local --primitive=both +python3 -m apache_beam.examples.inference.decision_models.fraud_review \ + --model=jev --primitive=both +``` + +A final three-element DirectRunner run returned these values. Beam may print +the rows in a different order: + +| message | Noul probability | Score | +| --- | ---: | ---: | +| billing address | 0.42 | 0.80 | +| duplicate charge | 0.32 | 0.85 | +| bypass verification | 0.97 | 3.00 | + +Only the bypass verification message crossed the `0.7` review threshold. + +## Tests + +Install Beam's test dependencies in the same virtual environment: + +```sh +python3 -m pip install -e 'sdks/python[test]' Review Comment: same adjustment to the python subdirectory here ########## sdks/python/apache_beam/examples/inference/decision_models/README.md: ########## @@ -0,0 +1,262 @@ +<!-- + Licensed to the Apache Software Foundation (ASF) under one + or more contributor license agreements. See the NOTICE file + distributed with this work for additional information + regarding copyright ownership. The ASF licenses this file + to you under the Apache License, Version 2.0 (the + "License"); you may not use this file except in compliance + with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, + software distributed under the License is distributed on an + "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + KIND, either express or implied. See the License for the + specific language governing permissions and limitations + under the License. +--> + +# Beam decision models + +These examples use the reusable decision primitives in +`apache_beam.ml.inference.decision`. A model evaluates named, typed questions +for each Beam element. `EvaluateDecisions` returns the original state, +normalized answers, provider metadata, and request latency. Policy and sinks +remain ordinary Beam transforms after evaluation. + +## Public API + +The core module contains `ChoiceQuestion`, `BooleanQuestion`, and +`ScoreQuestion`, their typed answers (`ChoiceAnswer`, `BooleanAnswer`, and +`ScoreAnswer`), the `DecisionModel` ABC, and the `EvaluateDecisions` +`PTransform`. + +| Question | Normalized answer | +| --- | --- | +| `ChoiceQuestion(instructions, criteria)` | `ChoiceAnswer.choice`, with optional label probabilities and confidence | +| `BooleanQuestion(instructions)` | `BooleanAnswer.probability` in `[0, 1]` | +| `ScoreQuestion(instructions, criteria)` | `ScoreAnswer.score` and `ScoreAnswer.probabilities`, keyed by zero-based integer level, with optional confidence | + +`JevDecisionModel` maps `BooleanQuestion` to Jev's `Noul`, returning +`BooleanAnswer.probability`; a missing confidence is `None`. + +`DecisionModel` is the adapter ABC. A subclass implements +`evaluate(state, questions)` and returns one matching typed answer per named +question in a `DecisionResponse`. Override `__enter__` and `__exit__` to reuse +a client, connection pool, or model weights. The transform enters the model +during DoFn setup and closes it during teardown for each DoFn instance. + +`EvaluateDecisions(model, questions)` is a `PTransform` from input states to +`DecisionResult`. Evaluation is synchronous, timestamps and windows are +preserved, exceptions propagate to the runner, and request, failure, and +latency metrics are recorded. + +## Architecture + +```mermaid +flowchart LR + input["PCollection[state]"] --> evaluate["EvaluateDecisions(model, questions)<br/>PTransform"] + model["DecisionModel adapter"] --> evaluate + evaluate --> result["PCollection[DecisionResult]"] + result --> policy["route_message<br/>policy and sink selection"] + policy --> sink["JSONL files<br/>or BigQuery tables"] +``` + + + +The optional Jev adapter lives in +`apache_beam.ml.inference.typesafe_inference`; the deterministic +`LocalDecisionModel` stays in +`apache_beam.examples.inference.decision_models.local_model`. Both implement +the same core contract. + +## Integration + +Questions and events are regular Python values. The model is the only +provider-specific argument to `EvaluateDecisions`: + +```python +import apache_beam as beam + +from apache_beam.ml.inference.decision import ChoiceQuestion +from apache_beam.ml.inference.decision import EvaluateDecisions +from apache_beam.examples.inference.decision_models.local_model import ( + LocalDecisionModel, +) +from apache_beam.ml.inference.typesafe_inference import JevDecisionModel + +questions = { + 'team': ChoiceQuestion('Which team should handle this message?', { + 'billing': 'Invoices', 'technical': 'Technical problems'}), +} + +def select_team(result): + return result.response.answers['team'].choice + +model = LocalDecisionModel() +# model = JevDecisionModel() + +with beam.Pipeline() as pipeline: + events = pipeline | 'Events' >> beam.Create( + [{'message': 'The API returns 503 after I rotated my key.'}]) + decisions = events | 'Evaluate decisions' >> EvaluateDecisions( + model, questions) + rows = decisions | 'Apply downstream policy' >> beam.Map(select_team) +``` + +The downstream policy stays unchanged when swapping adapters. The Jev client +is optional; install the existing example requirements and set +`TYPESAFE_API_KEY` when running it. + +## Install + +From a Beam checkout, activate a virtual environment and install the core SDK: + +```sh +python3 -m pip install -e sdks/python Review Comment: prefer to include a `cd sdks/python` step here to match other instructions about installing Python requirements from the python subdirectory -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
