Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
7 changes: 4 additions & 3 deletions src/sagemaker/transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -337,6 +337,7 @@ def transform_with_monitoring(
wait: bool = True,
pipeline_name: str = None,
role: str = None,
fail_on_violation: bool = True,
):
"""Runs a transform job with monitoring job.

Expand All@@ -352,7 +353,6 @@ def transform_with_monitoring(
]): the monitoring configuration used for run model monitoring.
monitoring_resource_config (`sagemaker.workflow.check_job_config.CheckJobConfig`):
the check job (processing job) cluster resource configuration.
transform_step_args (_JobStepArguments): the transform step transform arguments.
data (str): Input data location in S3 for the transform job
data_type (str): What the S3 location defines (default: 'S3Prefix').
Valid values:
Expand DownExpand Up@@ -400,8 +400,6 @@ def transform_with_monitoring(
monitor_before_transform (bgool): If to run data quality
or model explainability monitoring type,
a true value of this flag indicates running the check step before the transform job.
fail_on_violation (Union[bool, PipelineVariable]): A opt-out flag to not to fail the
check step when a violation is detected.
supplied_baseline_statistics (Union[str, PipelineVariable]): The S3 path
to the supplied statistics object representing the statistics JSON file
which will be used for drift to check (default: None).
Expand All@@ -411,6 +409,8 @@ def transform_with_monitoring(
wait (bool): To determine if needed to wait for the pipeline execution to complete
pipeline_name (str): The name of the Pipeline for the monitoring and transfrom step
role (str): Execution role
fail_on_violation (Union[bool, PipelineVariable]): A opt-out flag to not to fail the
check step when a violation is detected.
"""

transformer = self
Expand DownExpand Up@@ -454,6 +454,7 @@ def transform_with_monitoring(
monitor_before_transform=monitor_before_transform,
supplied_baseline_constraints=supplied_baseline_constraints,
supplied_baseline_statistics=supplied_baseline_statistics,
fail_on_violation=fail_on_violation,
Comment thread
keshav-chandak marked this conversation as resolved.
)

pipeline_name = (
Expand Down
64 changes: 64 additions & 0 deletions tests/integ/test_transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -709,3 +709,67 @@ def test_transformer_and_monitoring_job(
assert execution_step["StepStatus"] == "Succeeded"

xgb_model.delete_model()


def test_transformer_and_monitoring_job_to_pass_with_no_failure_in_violation(
pipeline_session,
sagemaker_session,
role,
pipeline_name,
check_job_config,
data_bias_check_config,
):
xgb_model_data_s3 = pipeline_session.upload_data(
path=os.path.join(os.path.join(DATA_DIR, "xgboost_abalone"), "xgb_model.tar.gz"),
key_prefix="integ-test-data/xgboost/model",
)
data_bias_supplied_baseline_constraints = Constraints.from_file_path(
constraints_file_path=os.path.join(
DATA_DIR, "pipeline/clarify_check_step/data_bias/bad_cases/analysis.json"
),
sagemaker_session=sagemaker_session,
).file_s3_uri

xgb_model = XGBoostModel(
model_data=xgb_model_data_s3,
framework_version="1.3-1",
role=role,
sagemaker_session=sagemaker_session,
entry_point=os.path.join(os.path.join(DATA_DIR, "xgboost_abalone"), "inference.py"),
enable_network_isolation=True,
)

xgb_model.deploy(_INSTANCE_COUNT, _INSTANCE_TYPE)

transform_output = f"s3://{sagemaker_session.default_bucket()}/{pipeline_name}Transform"
transformer = Transformer(
model_name=xgb_model.name,
strategy="SingleRecord",
instance_type="ml.m5.xlarge",
instance_count=1,
output_path=transform_output,
sagemaker_session=pipeline_session,
)

transform_input = pipeline_session.upload_data(
path=os.path.join(DATA_DIR, "xgboost_abalone", "abalone"),
key_prefix="integ-test-data/xgboost_abalone/abalone",
)

execution = transformer.transform_with_monitoring(
monitoring_config=data_bias_check_config,
monitoring_resource_config=check_job_config,
data=transform_input,
content_type="text/libsvm",
supplied_baseline_constraints=data_bias_supplied_baseline_constraints,
role=role,
fail_on_violation=False,
)

execution_steps = execution.list_steps()
assert len(execution_steps) == 2

for execution_step in execution_steps:
assert execution_step["StepStatus"] == "Succeeded"

xgb_model.delete_model()
61 changes: 60 additions & 1 deletion tests/unit/test_transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,6 +23,13 @@

from tests.integ import test_local_mode
from tests.unit import SAGEMAKER_CONFIG_TRANSFORM_JOB
from sagemaker.model_monitor import DatasetFormat
from sagemaker.workflow.quality_check_step import (
ModelQualityCheckConfig,
)
from sagemaker.workflow.check_job_config import CheckJobConfig

_CHECK_JOB_PREFIX = "CheckJobPrefix"

ROLE = "DummyRole"
REGION = "us-west-2"
Expand All@@ -49,6 +56,16 @@
"base_transform_job_name": JOB_NAME,
}

PROCESS_REQUEST_ARGS = {
"inputs": "processing_inputs",
"output_config": "output_config",
"job_name": "job_name",
"resources": "resource_config",
"stopping_condition": {"MaxRuntimeInSeconds": 3600},
"app_specification": "app_specification",
"experiment_config": {"ExperimentName": "AnExperiment"},
}

MODEL_DESC_PRIMARY_CONTAINER = {"PrimaryContainer": {"Image": IMAGE_URI}}

MODEL_DESC_CONTAINERS_ONLY = {"Containers": [{"Image": IMAGE_URI}]}
Expand All@@ -72,7 +89,7 @@ def mock_create_tar_file():

@pytest.fixture()
def sagemaker_session():
boto_mock = Mock(name="boto_session")
boto_mock = Mock(name="boto_session", region_name=REGION)
session = Mock(
name="sagemaker_session",
boto_session=boto_mock,
Expand DownExpand Up@@ -764,6 +781,48 @@ def test_stop_transform_job(sagemaker_session, transformer):
sagemaker_session.stop_transform_job.assert_called_once_with(name=JOB_NAME)


@patch("sagemaker.transformer.Transformer._retrieve_image_uri", return_value=IMAGE_URI)
@patch("sagemaker.workflow.pipeline.Pipeline.upsert", return_value={})
@patch("sagemaker.workflow.pipeline.Pipeline.start", return_value=Mock())
def test_transform_with_monitoring_create_and_starts_pipeline(
pipeline_start, upsert, image_uri, sagemaker_session, transformer
):

config = CheckJobConfig(
role=ROLE,
instance_count=1,
instance_type="ml.m5.xlarge",
volume_size_in_gb=60,
max_runtime_in_seconds=1800,
sagemaker_session=sagemaker_session,
base_job_name=_CHECK_JOB_PREFIX,
)

quality_check_config = ModelQualityCheckConfig(
baseline_dataset="s3://baseline_dataset_s3_url",
dataset_format=DatasetFormat.csv(header=True),
problem_type="BinaryClassification",
inference_attribute="quality_cfg_attr_value",
probability_attribute="quality_cfg_attr_value",
ground_truth_attribute="quality_cfg_attr_value",
probability_threshold_attribute="quality_cfg_attr_value",
post_analytics_processor_script="s3://my_bucket/data_quality/postprocessor.py",
output_s3_uri="s3://output_s3_uri",
)

transformer.transform_with_monitoring(
monitoring_config=quality_check_config,
monitoring_resource_config=config,
data=DATA,
content_type="text/libsvm",
supplied_baseline_constraints="supplied_baseline_constraints",
role=ROLE,
)

upsert.assert_called_once()
pipeline_start.assert_called_once()


def test_stop_transform_job_no_transform_job(transformer):
with pytest.raises(ValueError) as e:
transformer.stop_transform_job()
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
7 changes: 4 additions & 3 deletions src/sagemaker/transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -337,6 +337,7 @@ def transform_with_monitoring(
wait: bool = True,
pipeline_name: str = None,
role: str = None,
fail_on_violation: bool = True,
):
"""Runs a transform job with monitoring job.

Expand All@@ -352,7 +353,6 @@ def transform_with_monitoring(
]): the monitoring configuration used for run model monitoring.
monitoring_resource_config (`sagemaker.workflow.check_job_config.CheckJobConfig`):
the check job (processing job) cluster resource configuration.
transform_step_args (_JobStepArguments): the transform step transform arguments.
data (str): Input data location in S3 for the transform job
data_type (str): What the S3 location defines (default: 'S3Prefix').
Valid values:
Expand DownExpand Up@@ -400,8 +400,6 @@ def transform_with_monitoring(
monitor_before_transform (bgool): If to run data quality
or model explainability monitoring type,
a true value of this flag indicates running the check step before the transform job.
fail_on_violation (Union[bool, PipelineVariable]): A opt-out flag to not to fail the
check step when a violation is detected.
supplied_baseline_statistics (Union[str, PipelineVariable]): The S3 path
to the supplied statistics object representing the statistics JSON file
which will be used for drift to check (default: None).
Expand All@@ -411,6 +409,8 @@ def transform_with_monitoring(
wait (bool): To determine if needed to wait for the pipeline execution to complete
pipeline_name (str): The name of the Pipeline for the monitoring and transfrom step
role (str): Execution role
fail_on_violation (Union[bool, PipelineVariable]): A opt-out flag to not to fail the
check step when a violation is detected.
"""

transformer = self
Expand DownExpand Up@@ -454,6 +454,7 @@ def transform_with_monitoring(
monitor_before_transform=monitor_before_transform,
supplied_baseline_constraints=supplied_baseline_constraints,
supplied_baseline_statistics=supplied_baseline_statistics,
fail_on_violation=fail_on_violation,
Comment thread
keshav-chandak marked this conversation as resolved.
)

pipeline_name = (
Expand Down
64 changes: 64 additions & 0 deletions tests/integ/test_transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -709,3 +709,67 @@ def test_transformer_and_monitoring_job(
assert execution_step["StepStatus"] == "Succeeded"

xgb_model.delete_model()


def test_transformer_and_monitoring_job_to_pass_with_no_failure_in_violation(
pipeline_session,
sagemaker_session,
role,
pipeline_name,
check_job_config,
data_bias_check_config,
):
xgb_model_data_s3 = pipeline_session.upload_data(
path=os.path.join(os.path.join(DATA_DIR, "xgboost_abalone"), "xgb_model.tar.gz"),
key_prefix="integ-test-data/xgboost/model",
)
data_bias_supplied_baseline_constraints = Constraints.from_file_path(
constraints_file_path=os.path.join(
DATA_DIR, "pipeline/clarify_check_step/data_bias/bad_cases/analysis.json"
),
sagemaker_session=sagemaker_session,
).file_s3_uri

xgb_model = XGBoostModel(
model_data=xgb_model_data_s3,
framework_version="1.3-1",
role=role,
sagemaker_session=sagemaker_session,
entry_point=os.path.join(os.path.join(DATA_DIR, "xgboost_abalone"), "inference.py"),
enable_network_isolation=True,
)

xgb_model.deploy(_INSTANCE_COUNT, _INSTANCE_TYPE)

transform_output = f"s3://{sagemaker_session.default_bucket()}/{pipeline_name}Transform"
transformer = Transformer(
model_name=xgb_model.name,
strategy="SingleRecord",
instance_type="ml.m5.xlarge",
instance_count=1,
output_path=transform_output,
sagemaker_session=pipeline_session,
)

transform_input = pipeline_session.upload_data(
path=os.path.join(DATA_DIR, "xgboost_abalone", "abalone"),
key_prefix="integ-test-data/xgboost_abalone/abalone",
)

execution = transformer.transform_with_monitoring(
monitoring_config=data_bias_check_config,
monitoring_resource_config=check_job_config,
data=transform_input,
content_type="text/libsvm",
supplied_baseline_constraints=data_bias_supplied_baseline_constraints,
role=role,
fail_on_violation=False,
)

execution_steps = execution.list_steps()
assert len(execution_steps) == 2

for execution_step in execution_steps:
assert execution_step["StepStatus"] == "Succeeded"

xgb_model.delete_model()
61 changes: 60 additions & 1 deletion tests/unit/test_transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,6 +23,13 @@

from tests.integ import test_local_mode
from tests.unit import SAGEMAKER_CONFIG_TRANSFORM_JOB
from sagemaker.model_monitor import DatasetFormat
from sagemaker.workflow.quality_check_step import (
ModelQualityCheckConfig,
)
from sagemaker.workflow.check_job_config import CheckJobConfig

_CHECK_JOB_PREFIX = "CheckJobPrefix"

ROLE = "DummyRole"
REGION = "us-west-2"
Expand All@@ -49,6 +56,16 @@
"base_transform_job_name": JOB_NAME,
}

PROCESS_REQUEST_ARGS = {
"inputs": "processing_inputs",
"output_config": "output_config",
"job_name": "job_name",
"resources": "resource_config",
"stopping_condition": {"MaxRuntimeInSeconds": 3600},
"app_specification": "app_specification",
"experiment_config": {"ExperimentName": "AnExperiment"},
}

MODEL_DESC_PRIMARY_CONTAINER = {"PrimaryContainer": {"Image": IMAGE_URI}}

MODEL_DESC_CONTAINERS_ONLY = {"Containers": [{"Image": IMAGE_URI}]}
Expand All@@ -72,7 +89,7 @@ def mock_create_tar_file():

@pytest.fixture()
def sagemaker_session():
boto_mock = Mock(name="boto_session")
boto_mock = Mock(name="boto_session", region_name=REGION)
session = Mock(
name="sagemaker_session",
boto_session=boto_mock,
Expand DownExpand Up@@ -764,6 +781,48 @@ def test_stop_transform_job(sagemaker_session, transformer):
sagemaker_session.stop_transform_job.assert_called_once_with(name=JOB_NAME)


@patch("sagemaker.transformer.Transformer._retrieve_image_uri", return_value=IMAGE_URI)
@patch("sagemaker.workflow.pipeline.Pipeline.upsert", return_value={})
@patch("sagemaker.workflow.pipeline.Pipeline.start", return_value=Mock())
def test_transform_with_monitoring_create_and_starts_pipeline(
pipeline_start, upsert, image_uri, sagemaker_session, transformer
):

config = CheckJobConfig(
role=ROLE,
instance_count=1,
instance_type="ml.m5.xlarge",
volume_size_in_gb=60,
max_runtime_in_seconds=1800,
sagemaker_session=sagemaker_session,
base_job_name=_CHECK_JOB_PREFIX,
)

quality_check_config = ModelQualityCheckConfig(
baseline_dataset="s3://baseline_dataset_s3_url",
dataset_format=DatasetFormat.csv(header=True),
problem_type="BinaryClassification",
inference_attribute="quality_cfg_attr_value",
probability_attribute="quality_cfg_attr_value",
ground_truth_attribute="quality_cfg_attr_value",
probability_threshold_attribute="quality_cfg_attr_value",
post_analytics_processor_script="s3://my_bucket/data_quality/postprocessor.py",
output_s3_uri="s3://output_s3_uri",
)

transformer.transform_with_monitoring(
monitoring_config=quality_check_config,
monitoring_resource_config=config,
data=DATA,
content_type="text/libsvm",
supplied_baseline_constraints="supplied_baseline_constraints",
role=ROLE,
)

upsert.assert_called_once()
pipeline_start.assert_called_once()


def test_stop_transform_job_no_transform_job(transformer):
with pytest.raises(ValueError) as e:
transformer.stop_transform_job()
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
7 changes: 4 additions & 3 deletions src/sagemaker/transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -337,6 +337,7 @@ def transform_with_monitoring(
wait: bool = True,
pipeline_name: str = None,
role: str = None,
fail_on_violation: bool = True,
):
"""Runs a transform job with monitoring job.

Expand All@@ -352,7 +353,6 @@ def transform_with_monitoring(
]): the monitoring configuration used for run model monitoring.
monitoring_resource_config (`sagemaker.workflow.check_job_config.CheckJobConfig`):
the check job (processing job) cluster resource configuration.
transform_step_args (_JobStepArguments): the transform step transform arguments.
data (str): Input data location in S3 for the transform job
data_type (str): What the S3 location defines (default: 'S3Prefix').
Valid values:
Expand DownExpand Up@@ -400,8 +400,6 @@ def transform_with_monitoring(
monitor_before_transform (bgool): If to run data quality
or model explainability monitoring type,
a true value of this flag indicates running the check step before the transform job.
fail_on_violation (Union[bool, PipelineVariable]): A opt-out flag to not to fail the
check step when a violation is detected.
supplied_baseline_statistics (Union[str, PipelineVariable]): The S3 path
to the supplied statistics object representing the statistics JSON file
which will be used for drift to check (default: None).
Expand All@@ -411,6 +409,8 @@ def transform_with_monitoring(
wait (bool): To determine if needed to wait for the pipeline execution to complete
pipeline_name (str): The name of the Pipeline for the monitoring and transfrom step
role (str): Execution role
fail_on_violation (Union[bool, PipelineVariable]): A opt-out flag to not to fail the
check step when a violation is detected.
"""

transformer = self
Expand DownExpand Up@@ -454,6 +454,7 @@ def transform_with_monitoring(
monitor_before_transform=monitor_before_transform,
supplied_baseline_constraints=supplied_baseline_constraints,
supplied_baseline_statistics=supplied_baseline_statistics,
fail_on_violation=fail_on_violation,
Comment thread
keshav-chandak marked this conversation as resolved.
)

pipeline_name = (
Expand Down
64 changes: 64 additions & 0 deletions tests/integ/test_transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -709,3 +709,67 @@ def test_transformer_and_monitoring_job(
assert execution_step["StepStatus"] == "Succeeded"

xgb_model.delete_model()


def test_transformer_and_monitoring_job_to_pass_with_no_failure_in_violation(
pipeline_session,
sagemaker_session,
role,
pipeline_name,
check_job_config,
data_bias_check_config,
):
xgb_model_data_s3 = pipeline_session.upload_data(
path=os.path.join(os.path.join(DATA_DIR, "xgboost_abalone"), "xgb_model.tar.gz"),
key_prefix="integ-test-data/xgboost/model",
)
data_bias_supplied_baseline_constraints = Constraints.from_file_path(
constraints_file_path=os.path.join(
DATA_DIR, "pipeline/clarify_check_step/data_bias/bad_cases/analysis.json"
),
sagemaker_session=sagemaker_session,
).file_s3_uri

xgb_model = XGBoostModel(
model_data=xgb_model_data_s3,
framework_version="1.3-1",
role=role,
sagemaker_session=sagemaker_session,
entry_point=os.path.join(os.path.join(DATA_DIR, "xgboost_abalone"), "inference.py"),
enable_network_isolation=True,
)

xgb_model.deploy(_INSTANCE_COUNT, _INSTANCE_TYPE)

transform_output = f"s3://{sagemaker_session.default_bucket()}/{pipeline_name}Transform"
transformer = Transformer(
model_name=xgb_model.name,
strategy="SingleRecord",
instance_type="ml.m5.xlarge",
instance_count=1,
output_path=transform_output,
sagemaker_session=pipeline_session,
)

transform_input = pipeline_session.upload_data(
path=os.path.join(DATA_DIR, "xgboost_abalone", "abalone"),
key_prefix="integ-test-data/xgboost_abalone/abalone",
)

execution = transformer.transform_with_monitoring(
monitoring_config=data_bias_check_config,
monitoring_resource_config=check_job_config,
data=transform_input,
content_type="text/libsvm",
supplied_baseline_constraints=data_bias_supplied_baseline_constraints,
role=role,
fail_on_violation=False,
)

execution_steps = execution.list_steps()
assert len(execution_steps) == 2

for execution_step in execution_steps:
assert execution_step["StepStatus"] == "Succeeded"

xgb_model.delete_model()
61 changes: 60 additions & 1 deletion tests/unit/test_transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,6 +23,13 @@

from tests.integ import test_local_mode
from tests.unit import SAGEMAKER_CONFIG_TRANSFORM_JOB
from sagemaker.model_monitor import DatasetFormat
from sagemaker.workflow.quality_check_step import (
ModelQualityCheckConfig,
)
from sagemaker.workflow.check_job_config import CheckJobConfig

_CHECK_JOB_PREFIX = "CheckJobPrefix"

ROLE = "DummyRole"
REGION = "us-west-2"
Expand All@@ -49,6 +56,16 @@
"base_transform_job_name": JOB_NAME,
}

PROCESS_REQUEST_ARGS = {
"inputs": "processing_inputs",
"output_config": "output_config",
"job_name": "job_name",
"resources": "resource_config",
"stopping_condition": {"MaxRuntimeInSeconds": 3600},
"app_specification": "app_specification",
"experiment_config": {"ExperimentName": "AnExperiment"},
}

MODEL_DESC_PRIMARY_CONTAINER = {"PrimaryContainer": {"Image": IMAGE_URI}}

MODEL_DESC_CONTAINERS_ONLY = {"Containers": [{"Image": IMAGE_URI}]}
Expand All@@ -72,7 +89,7 @@ def mock_create_tar_file():

@pytest.fixture()
def sagemaker_session():
boto_mock = Mock(name="boto_session")
boto_mock = Mock(name="boto_session", region_name=REGION)
session = Mock(
name="sagemaker_session",
boto_session=boto_mock,
Expand DownExpand Up@@ -764,6 +781,48 @@ def test_stop_transform_job(sagemaker_session, transformer):
sagemaker_session.stop_transform_job.assert_called_once_with(name=JOB_NAME)


@patch("sagemaker.transformer.Transformer._retrieve_image_uri", return_value=IMAGE_URI)
@patch("sagemaker.workflow.pipeline.Pipeline.upsert", return_value={})
@patch("sagemaker.workflow.pipeline.Pipeline.start", return_value=Mock())
def test_transform_with_monitoring_create_and_starts_pipeline(
pipeline_start, upsert, image_uri, sagemaker_session, transformer
):

config = CheckJobConfig(
role=ROLE,
instance_count=1,
instance_type="ml.m5.xlarge",
volume_size_in_gb=60,
max_runtime_in_seconds=1800,
sagemaker_session=sagemaker_session,
base_job_name=_CHECK_JOB_PREFIX,
)

quality_check_config = ModelQualityCheckConfig(
baseline_dataset="s3://baseline_dataset_s3_url",
dataset_format=DatasetFormat.csv(header=True),
problem_type="BinaryClassification",
inference_attribute="quality_cfg_attr_value",
probability_attribute="quality_cfg_attr_value",
ground_truth_attribute="quality_cfg_attr_value",
probability_threshold_attribute="quality_cfg_attr_value",
post_analytics_processor_script="s3://my_bucket/data_quality/postprocessor.py",
output_s3_uri="s3://output_s3_uri",
)

transformer.transform_with_monitoring(
monitoring_config=quality_check_config,
monitoring_resource_config=config,
data=DATA,
content_type="text/libsvm",
supplied_baseline_constraints="supplied_baseline_constraints",
role=ROLE,
)

upsert.assert_called_once()
pipeline_start.assert_called_once()


def test_stop_transform_job_no_transform_job(transformer):
with pytest.raises(ValueError) as e:
transformer.stop_transform_job()
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
7 changes: 4 additions & 3 deletions src/sagemaker/transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -337,6 +337,7 @@ def transform_with_monitoring(
wait: bool = True,
pipeline_name: str = None,
role: str = None,
fail_on_violation: bool = True,
):
"""Runs a transform job with monitoring job.

Expand All@@ -352,7 +353,6 @@ def transform_with_monitoring(
]): the monitoring configuration used for run model monitoring.
monitoring_resource_config (`sagemaker.workflow.check_job_config.CheckJobConfig`):
the check job (processing job) cluster resource configuration.
transform_step_args (_JobStepArguments): the transform step transform arguments.
data (str): Input data location in S3 for the transform job
data_type (str): What the S3 location defines (default: 'S3Prefix').
Valid values:
Expand DownExpand Up@@ -400,8 +400,6 @@ def transform_with_monitoring(
monitor_before_transform (bgool): If to run data quality
or model explainability monitoring type,
a true value of this flag indicates running the check step before the transform job.
fail_on_violation (Union[bool, PipelineVariable]): A opt-out flag to not to fail the
check step when a violation is detected.
supplied_baseline_statistics (Union[str, PipelineVariable]): The S3 path
to the supplied statistics object representing the statistics JSON file
which will be used for drift to check (default: None).
Expand All@@ -411,6 +409,8 @@ def transform_with_monitoring(
wait (bool): To determine if needed to wait for the pipeline execution to complete
pipeline_name (str): The name of the Pipeline for the monitoring and transfrom step
role (str): Execution role
fail_on_violation (Union[bool, PipelineVariable]): A opt-out flag to not to fail the
check step when a violation is detected.
"""

transformer = self
Expand DownExpand Up@@ -454,6 +454,7 @@ def transform_with_monitoring(
monitor_before_transform=monitor_before_transform,
supplied_baseline_constraints=supplied_baseline_constraints,
supplied_baseline_statistics=supplied_baseline_statistics,
fail_on_violation=fail_on_violation,
Comment thread
keshav-chandak marked this conversation as resolved.
)

pipeline_name = (
Expand Down
64 changes: 64 additions & 0 deletions tests/integ/test_transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -709,3 +709,67 @@ def test_transformer_and_monitoring_job(
assert execution_step["StepStatus"] == "Succeeded"

xgb_model.delete_model()


def test_transformer_and_monitoring_job_to_pass_with_no_failure_in_violation(
pipeline_session,
sagemaker_session,
role,
pipeline_name,
check_job_config,
data_bias_check_config,
):
xgb_model_data_s3 = pipeline_session.upload_data(
path=os.path.join(os.path.join(DATA_DIR, "xgboost_abalone"), "xgb_model.tar.gz"),
key_prefix="integ-test-data/xgboost/model",
)
data_bias_supplied_baseline_constraints = Constraints.from_file_path(
constraints_file_path=os.path.join(
DATA_DIR, "pipeline/clarify_check_step/data_bias/bad_cases/analysis.json"
),
sagemaker_session=sagemaker_session,
).file_s3_uri

xgb_model = XGBoostModel(
model_data=xgb_model_data_s3,
framework_version="1.3-1",
role=role,
sagemaker_session=sagemaker_session,
entry_point=os.path.join(os.path.join(DATA_DIR, "xgboost_abalone"), "inference.py"),
enable_network_isolation=True,
)

xgb_model.deploy(_INSTANCE_COUNT, _INSTANCE_TYPE)

transform_output = f"s3://{sagemaker_session.default_bucket()}/{pipeline_name}Transform"
transformer = Transformer(
model_name=xgb_model.name,
strategy="SingleRecord",
instance_type="ml.m5.xlarge",
instance_count=1,
output_path=transform_output,
sagemaker_session=pipeline_session,
)

transform_input = pipeline_session.upload_data(
path=os.path.join(DATA_DIR, "xgboost_abalone", "abalone"),
key_prefix="integ-test-data/xgboost_abalone/abalone",
)

execution = transformer.transform_with_monitoring(
monitoring_config=data_bias_check_config,
monitoring_resource_config=check_job_config,
data=transform_input,
content_type="text/libsvm",
supplied_baseline_constraints=data_bias_supplied_baseline_constraints,
role=role,
fail_on_violation=False,
)

execution_steps = execution.list_steps()
assert len(execution_steps) == 2

for execution_step in execution_steps:
assert execution_step["StepStatus"] == "Succeeded"

xgb_model.delete_model()
61 changes: 60 additions & 1 deletion tests/unit/test_transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,6 +23,13 @@

from tests.integ import test_local_mode
from tests.unit import SAGEMAKER_CONFIG_TRANSFORM_JOB
from sagemaker.model_monitor import DatasetFormat
from sagemaker.workflow.quality_check_step import (
ModelQualityCheckConfig,
)
from sagemaker.workflow.check_job_config import CheckJobConfig

_CHECK_JOB_PREFIX = "CheckJobPrefix"

ROLE = "DummyRole"
REGION = "us-west-2"
Expand All@@ -49,6 +56,16 @@
"base_transform_job_name": JOB_NAME,
}

PROCESS_REQUEST_ARGS = {
"inputs": "processing_inputs",
"output_config": "output_config",
"job_name": "job_name",
"resources": "resource_config",
"stopping_condition": {"MaxRuntimeInSeconds": 3600},
"app_specification": "app_specification",
"experiment_config": {"ExperimentName": "AnExperiment"},
}

MODEL_DESC_PRIMARY_CONTAINER = {"PrimaryContainer": {"Image": IMAGE_URI}}

MODEL_DESC_CONTAINERS_ONLY = {"Containers": [{"Image": IMAGE_URI}]}
Expand All@@ -72,7 +89,7 @@ def mock_create_tar_file():

@pytest.fixture()
def sagemaker_session():
boto_mock = Mock(name="boto_session")
boto_mock = Mock(name="boto_session", region_name=REGION)
session = Mock(
name="sagemaker_session",
boto_session=boto_mock,
Expand DownExpand Up@@ -764,6 +781,48 @@ def test_stop_transform_job(sagemaker_session, transformer):
sagemaker_session.stop_transform_job.assert_called_once_with(name=JOB_NAME)


@patch("sagemaker.transformer.Transformer._retrieve_image_uri", return_value=IMAGE_URI)
@patch("sagemaker.workflow.pipeline.Pipeline.upsert", return_value={})
@patch("sagemaker.workflow.pipeline.Pipeline.start", return_value=Mock())
def test_transform_with_monitoring_create_and_starts_pipeline(
pipeline_start, upsert, image_uri, sagemaker_session, transformer
):

config = CheckJobConfig(
role=ROLE,
instance_count=1,
instance_type="ml.m5.xlarge",
volume_size_in_gb=60,
max_runtime_in_seconds=1800,
sagemaker_session=sagemaker_session,
base_job_name=_CHECK_JOB_PREFIX,
)

quality_check_config = ModelQualityCheckConfig(
baseline_dataset="s3://baseline_dataset_s3_url",
dataset_format=DatasetFormat.csv(header=True),
problem_type="BinaryClassification",
inference_attribute="quality_cfg_attr_value",
probability_attribute="quality_cfg_attr_value",
ground_truth_attribute="quality_cfg_attr_value",
probability_threshold_attribute="quality_cfg_attr_value",
post_analytics_processor_script="s3://my_bucket/data_quality/postprocessor.py",
output_s3_uri="s3://output_s3_uri",
)

transformer.transform_with_monitoring(
monitoring_config=quality_check_config,
monitoring_resource_config=config,
data=DATA,
content_type="text/libsvm",
supplied_baseline_constraints="supplied_baseline_constraints",
role=ROLE,
)

upsert.assert_called_once()
pipeline_start.assert_called_once()


def test_stop_transform_job_no_transform_job(transformer):
with pytest.raises(ValueError) as e:
transformer.stop_transform_job()
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
7 changes: 4 additions & 3 deletions src/sagemaker/transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -337,6 +337,7 @@ def transform_with_monitoring(
wait: bool = True,
pipeline_name: str = None,
role: str = None,
fail_on_violation: bool = True,
):
"""Runs a transform job with monitoring job.

Expand All@@ -352,7 +353,6 @@ def transform_with_monitoring(
]): the monitoring configuration used for run model monitoring.
monitoring_resource_config (`sagemaker.workflow.check_job_config.CheckJobConfig`):
the check job (processing job) cluster resource configuration.
transform_step_args (_JobStepArguments): the transform step transform arguments.
data (str): Input data location in S3 for the transform job
data_type (str): What the S3 location defines (default: 'S3Prefix').
Valid values:
Expand DownExpand Up@@ -400,8 +400,6 @@ def transform_with_monitoring(
monitor_before_transform (bgool): If to run data quality
or model explainability monitoring type,
a true value of this flag indicates running the check step before the transform job.
fail_on_violation (Union[bool, PipelineVariable]): A opt-out flag to not to fail the
check step when a violation is detected.
supplied_baseline_statistics (Union[str, PipelineVariable]): The S3 path
to the supplied statistics object representing the statistics JSON file
which will be used for drift to check (default: None).
Expand All@@ -411,6 +409,8 @@ def transform_with_monitoring(
wait (bool): To determine if needed to wait for the pipeline execution to complete
pipeline_name (str): The name of the Pipeline for the monitoring and transfrom step
role (str): Execution role
fail_on_violation (Union[bool, PipelineVariable]): A opt-out flag to not to fail the
check step when a violation is detected.
"""

transformer = self
Expand DownExpand Up@@ -454,6 +454,7 @@ def transform_with_monitoring(
monitor_before_transform=monitor_before_transform,
supplied_baseline_constraints=supplied_baseline_constraints,
supplied_baseline_statistics=supplied_baseline_statistics,
fail_on_violation=fail_on_violation,
Comment thread
keshav-chandak marked this conversation as resolved.
)

pipeline_name = (
Expand Down
64 changes: 64 additions & 0 deletions tests/integ/test_transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -709,3 +709,67 @@ def test_transformer_and_monitoring_job(
assert execution_step["StepStatus"] == "Succeeded"

xgb_model.delete_model()


def test_transformer_and_monitoring_job_to_pass_with_no_failure_in_violation(
pipeline_session,
sagemaker_session,
role,
pipeline_name,
check_job_config,
data_bias_check_config,
):
xgb_model_data_s3 = pipeline_session.upload_data(
path=os.path.join(os.path.join(DATA_DIR, "xgboost_abalone"), "xgb_model.tar.gz"),
key_prefix="integ-test-data/xgboost/model",
)
data_bias_supplied_baseline_constraints = Constraints.from_file_path(
constraints_file_path=os.path.join(
DATA_DIR, "pipeline/clarify_check_step/data_bias/bad_cases/analysis.json"
),
sagemaker_session=sagemaker_session,
).file_s3_uri

xgb_model = XGBoostModel(
model_data=xgb_model_data_s3,
framework_version="1.3-1",
role=role,
sagemaker_session=sagemaker_session,
entry_point=os.path.join(os.path.join(DATA_DIR, "xgboost_abalone"), "inference.py"),
enable_network_isolation=True,
)

xgb_model.deploy(_INSTANCE_COUNT, _INSTANCE_TYPE)

transform_output = f"s3://{sagemaker_session.default_bucket()}/{pipeline_name}Transform"
transformer = Transformer(
model_name=xgb_model.name,
strategy="SingleRecord",
instance_type="ml.m5.xlarge",
instance_count=1,
output_path=transform_output,
sagemaker_session=pipeline_session,
)

transform_input = pipeline_session.upload_data(
path=os.path.join(DATA_DIR, "xgboost_abalone", "abalone"),
key_prefix="integ-test-data/xgboost_abalone/abalone",
)

execution = transformer.transform_with_monitoring(
monitoring_config=data_bias_check_config,
monitoring_resource_config=check_job_config,
data=transform_input,
content_type="text/libsvm",
supplied_baseline_constraints=data_bias_supplied_baseline_constraints,
role=role,
fail_on_violation=False,
)

execution_steps = execution.list_steps()
assert len(execution_steps) == 2

for execution_step in execution_steps:
assert execution_step["StepStatus"] == "Succeeded"

xgb_model.delete_model()
61 changes: 60 additions & 1 deletion tests/unit/test_transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,6 +23,13 @@

from tests.integ import test_local_mode
from tests.unit import SAGEMAKER_CONFIG_TRANSFORM_JOB
from sagemaker.model_monitor import DatasetFormat
from sagemaker.workflow.quality_check_step import (
ModelQualityCheckConfig,
)
from sagemaker.workflow.check_job_config import CheckJobConfig

_CHECK_JOB_PREFIX = "CheckJobPrefix"

ROLE = "DummyRole"
REGION = "us-west-2"
Expand All@@ -49,6 +56,16 @@
"base_transform_job_name": JOB_NAME,
}

PROCESS_REQUEST_ARGS = {
"inputs": "processing_inputs",
"output_config": "output_config",
"job_name": "job_name",
"resources": "resource_config",
"stopping_condition": {"MaxRuntimeInSeconds": 3600},
"app_specification": "app_specification",
"experiment_config": {"ExperimentName": "AnExperiment"},
}

MODEL_DESC_PRIMARY_CONTAINER = {"PrimaryContainer": {"Image": IMAGE_URI}}

MODEL_DESC_CONTAINERS_ONLY = {"Containers": [{"Image": IMAGE_URI}]}
Expand All@@ -72,7 +89,7 @@ def mock_create_tar_file():

@pytest.fixture()
def sagemaker_session():
boto_mock = Mock(name="boto_session")
boto_mock = Mock(name="boto_session", region_name=REGION)
session = Mock(
name="sagemaker_session",
boto_session=boto_mock,
Expand DownExpand Up@@ -764,6 +781,48 @@ def test_stop_transform_job(sagemaker_session, transformer):
sagemaker_session.stop_transform_job.assert_called_once_with(name=JOB_NAME)


@patch("sagemaker.transformer.Transformer._retrieve_image_uri", return_value=IMAGE_URI)
@patch("sagemaker.workflow.pipeline.Pipeline.upsert", return_value={})
@patch("sagemaker.workflow.pipeline.Pipeline.start", return_value=Mock())
def test_transform_with_monitoring_create_and_starts_pipeline(
pipeline_start, upsert, image_uri, sagemaker_session, transformer
):

config = CheckJobConfig(
role=ROLE,
instance_count=1,
instance_type="ml.m5.xlarge",
volume_size_in_gb=60,
max_runtime_in_seconds=1800,
sagemaker_session=sagemaker_session,
base_job_name=_CHECK_JOB_PREFIX,
)

quality_check_config = ModelQualityCheckConfig(
baseline_dataset="s3://baseline_dataset_s3_url",
dataset_format=DatasetFormat.csv(header=True),
problem_type="BinaryClassification",
inference_attribute="quality_cfg_attr_value",
probability_attribute="quality_cfg_attr_value",
ground_truth_attribute="quality_cfg_attr_value",
probability_threshold_attribute="quality_cfg_attr_value",
post_analytics_processor_script="s3://my_bucket/data_quality/postprocessor.py",
output_s3_uri="s3://output_s3_uri",
)

transformer.transform_with_monitoring(
monitoring_config=quality_check_config,
monitoring_resource_config=config,
data=DATA,
content_type="text/libsvm",
supplied_baseline_constraints="supplied_baseline_constraints",
role=ROLE,
)

upsert.assert_called_once()
pipeline_start.assert_called_once()


def test_stop_transform_job_no_transform_job(transformer):
with pytest.raises(ValueError) as e:
transformer.stop_transform_job()
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
7 changes: 4 additions & 3 deletions src/sagemaker/transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -337,6 +337,7 @@ def transform_with_monitoring(
wait: bool = True,
pipeline_name: str = None,
role: str = None,
fail_on_violation: bool = True,
):
"""Runs a transform job with monitoring job.

Expand All@@ -352,7 +353,6 @@ def transform_with_monitoring(
]): the monitoring configuration used for run model monitoring.
monitoring_resource_config (`sagemaker.workflow.check_job_config.CheckJobConfig`):
the check job (processing job) cluster resource configuration.
transform_step_args (_JobStepArguments): the transform step transform arguments.
data (str): Input data location in S3 for the transform job
data_type (str): What the S3 location defines (default: 'S3Prefix').
Valid values:
Expand DownExpand Up@@ -400,8 +400,6 @@ def transform_with_monitoring(
monitor_before_transform (bgool): If to run data quality
or model explainability monitoring type,
a true value of this flag indicates running the check step before the transform job.
fail_on_violation (Union[bool, PipelineVariable]): A opt-out flag to not to fail the
check step when a violation is detected.
supplied_baseline_statistics (Union[str, PipelineVariable]): The S3 path
to the supplied statistics object representing the statistics JSON file
which will be used for drift to check (default: None).
Expand All@@ -411,6 +409,8 @@ def transform_with_monitoring(
wait (bool): To determine if needed to wait for the pipeline execution to complete
pipeline_name (str): The name of the Pipeline for the monitoring and transfrom step
role (str): Execution role
fail_on_violation (Union[bool, PipelineVariable]): A opt-out flag to not to fail the
check step when a violation is detected.
"""

transformer = self
Expand DownExpand Up@@ -454,6 +454,7 @@ def transform_with_monitoring(
monitor_before_transform=monitor_before_transform,
supplied_baseline_constraints=supplied_baseline_constraints,
supplied_baseline_statistics=supplied_baseline_statistics,
fail_on_violation=fail_on_violation,
Comment thread
keshav-chandak marked this conversation as resolved.
)

pipeline_name = (
Expand Down
64 changes: 64 additions & 0 deletions tests/integ/test_transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -709,3 +709,67 @@ def test_transformer_and_monitoring_job(
assert execution_step["StepStatus"] == "Succeeded"

xgb_model.delete_model()


def test_transformer_and_monitoring_job_to_pass_with_no_failure_in_violation(
pipeline_session,
sagemaker_session,
role,
pipeline_name,
check_job_config,
data_bias_check_config,
):
xgb_model_data_s3 = pipeline_session.upload_data(
path=os.path.join(os.path.join(DATA_DIR, "xgboost_abalone"), "xgb_model.tar.gz"),
key_prefix="integ-test-data/xgboost/model",
)
data_bias_supplied_baseline_constraints = Constraints.from_file_path(
constraints_file_path=os.path.join(
DATA_DIR, "pipeline/clarify_check_step/data_bias/bad_cases/analysis.json"
),
sagemaker_session=sagemaker_session,
).file_s3_uri

xgb_model = XGBoostModel(
model_data=xgb_model_data_s3,
framework_version="1.3-1",
role=role,
sagemaker_session=sagemaker_session,
entry_point=os.path.join(os.path.join(DATA_DIR, "xgboost_abalone"), "inference.py"),
enable_network_isolation=True,
)

xgb_model.deploy(_INSTANCE_COUNT, _INSTANCE_TYPE)

transform_output = f"s3://{sagemaker_session.default_bucket()}/{pipeline_name}Transform"
transformer = Transformer(
model_name=xgb_model.name,
strategy="SingleRecord",
instance_type="ml.m5.xlarge",
instance_count=1,
output_path=transform_output,
sagemaker_session=pipeline_session,
)

transform_input = pipeline_session.upload_data(
path=os.path.join(DATA_DIR, "xgboost_abalone", "abalone"),
key_prefix="integ-test-data/xgboost_abalone/abalone",
)

execution = transformer.transform_with_monitoring(
monitoring_config=data_bias_check_config,
monitoring_resource_config=check_job_config,
data=transform_input,
content_type="text/libsvm",
supplied_baseline_constraints=data_bias_supplied_baseline_constraints,
role=role,
fail_on_violation=False,
)

execution_steps = execution.list_steps()
assert len(execution_steps) == 2

for execution_step in execution_steps:
assert execution_step["StepStatus"] == "Succeeded"

xgb_model.delete_model()
61 changes: 60 additions & 1 deletion tests/unit/test_transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,6 +23,13 @@

from tests.integ import test_local_mode
from tests.unit import SAGEMAKER_CONFIG_TRANSFORM_JOB
from sagemaker.model_monitor import DatasetFormat
from sagemaker.workflow.quality_check_step import (
ModelQualityCheckConfig,
)
from sagemaker.workflow.check_job_config import CheckJobConfig

_CHECK_JOB_PREFIX = "CheckJobPrefix"

ROLE = "DummyRole"
REGION = "us-west-2"
Expand All@@ -49,6 +56,16 @@
"base_transform_job_name": JOB_NAME,
}

PROCESS_REQUEST_ARGS = {
"inputs": "processing_inputs",
"output_config": "output_config",
"job_name": "job_name",
"resources": "resource_config",
"stopping_condition": {"MaxRuntimeInSeconds": 3600},
"app_specification": "app_specification",
"experiment_config": {"ExperimentName": "AnExperiment"},
}

MODEL_DESC_PRIMARY_CONTAINER = {"PrimaryContainer": {"Image": IMAGE_URI}}

MODEL_DESC_CONTAINERS_ONLY = {"Containers": [{"Image": IMAGE_URI}]}
Expand All@@ -72,7 +89,7 @@ def mock_create_tar_file():

@pytest.fixture()
def sagemaker_session():
boto_mock = Mock(name="boto_session")
boto_mock = Mock(name="boto_session", region_name=REGION)
session = Mock(
name="sagemaker_session",
boto_session=boto_mock,
Expand DownExpand Up@@ -764,6 +781,48 @@ def test_stop_transform_job(sagemaker_session, transformer):
sagemaker_session.stop_transform_job.assert_called_once_with(name=JOB_NAME)


@patch("sagemaker.transformer.Transformer._retrieve_image_uri", return_value=IMAGE_URI)
@patch("sagemaker.workflow.pipeline.Pipeline.upsert", return_value={})
@patch("sagemaker.workflow.pipeline.Pipeline.start", return_value=Mock())
def test_transform_with_monitoring_create_and_starts_pipeline(
pipeline_start, upsert, image_uri, sagemaker_session, transformer
):

config = CheckJobConfig(
role=ROLE,
instance_count=1,
instance_type="ml.m5.xlarge",
volume_size_in_gb=60,
max_runtime_in_seconds=1800,
sagemaker_session=sagemaker_session,
base_job_name=_CHECK_JOB_PREFIX,
)

quality_check_config = ModelQualityCheckConfig(
baseline_dataset="s3://baseline_dataset_s3_url",
dataset_format=DatasetFormat.csv(header=True),
problem_type="BinaryClassification",
inference_attribute="quality_cfg_attr_value",
probability_attribute="quality_cfg_attr_value",
ground_truth_attribute="quality_cfg_attr_value",
probability_threshold_attribute="quality_cfg_attr_value",
post_analytics_processor_script="s3://my_bucket/data_quality/postprocessor.py",
output_s3_uri="s3://output_s3_uri",
)

transformer.transform_with_monitoring(
monitoring_config=quality_check_config,
monitoring_resource_config=config,
data=DATA,
content_type="text/libsvm",
supplied_baseline_constraints="supplied_baseline_constraints",
role=ROLE,
)

upsert.assert_called_once()
pipeline_start.assert_called_once()


def test_stop_transform_job_no_transform_job(transformer):
with pytest.raises(ValueError) as e:
transformer.stop_transform_job()
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
7 changes: 4 additions & 3 deletions src/sagemaker/transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -337,6 +337,7 @@ def transform_with_monitoring(
wait: bool = True,
pipeline_name: str = None,
role: str = None,
fail_on_violation: bool = True,
):
"""Runs a transform job with monitoring job.

Expand All@@ -352,7 +353,6 @@ def transform_with_monitoring(
]): the monitoring configuration used for run model monitoring.
monitoring_resource_config (`sagemaker.workflow.check_job_config.CheckJobConfig`):
the check job (processing job) cluster resource configuration.
transform_step_args (_JobStepArguments): the transform step transform arguments.
data (str): Input data location in S3 for the transform job
data_type (str): What the S3 location defines (default: 'S3Prefix').
Valid values:
Expand DownExpand Up@@ -400,8 +400,6 @@ def transform_with_monitoring(
monitor_before_transform (bgool): If to run data quality
or model explainability monitoring type,
a true value of this flag indicates running the check step before the transform job.
fail_on_violation (Union[bool, PipelineVariable]): A opt-out flag to not to fail the
check step when a violation is detected.
supplied_baseline_statistics (Union[str, PipelineVariable]): The S3 path
to the supplied statistics object representing the statistics JSON file
which will be used for drift to check (default: None).
Expand All@@ -411,6 +409,8 @@ def transform_with_monitoring(
wait (bool): To determine if needed to wait for the pipeline execution to complete
pipeline_name (str): The name of the Pipeline for the monitoring and transfrom step
role (str): Execution role
fail_on_violation (Union[bool, PipelineVariable]): A opt-out flag to not to fail the
check step when a violation is detected.
"""

transformer = self
Expand DownExpand Up@@ -454,6 +454,7 @@ def transform_with_monitoring(
monitor_before_transform=monitor_before_transform,
supplied_baseline_constraints=supplied_baseline_constraints,
supplied_baseline_statistics=supplied_baseline_statistics,
fail_on_violation=fail_on_violation,
Comment thread
keshav-chandak marked this conversation as resolved.
)

pipeline_name = (
Expand Down
64 changes: 64 additions & 0 deletions tests/integ/test_transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -709,3 +709,67 @@ def test_transformer_and_monitoring_job(
assert execution_step["StepStatus"] == "Succeeded"

xgb_model.delete_model()


def test_transformer_and_monitoring_job_to_pass_with_no_failure_in_violation(
pipeline_session,
sagemaker_session,
role,
pipeline_name,
check_job_config,
data_bias_check_config,
):
xgb_model_data_s3 = pipeline_session.upload_data(
path=os.path.join(os.path.join(DATA_DIR, "xgboost_abalone"), "xgb_model.tar.gz"),
key_prefix="integ-test-data/xgboost/model",
)
data_bias_supplied_baseline_constraints = Constraints.from_file_path(
constraints_file_path=os.path.join(
DATA_DIR, "pipeline/clarify_check_step/data_bias/bad_cases/analysis.json"
),
sagemaker_session=sagemaker_session,
).file_s3_uri

xgb_model = XGBoostModel(
model_data=xgb_model_data_s3,
framework_version="1.3-1",
role=role,
sagemaker_session=sagemaker_session,
entry_point=os.path.join(os.path.join(DATA_DIR, "xgboost_abalone"), "inference.py"),
enable_network_isolation=True,
)

xgb_model.deploy(_INSTANCE_COUNT, _INSTANCE_TYPE)

transform_output = f"s3://{sagemaker_session.default_bucket()}/{pipeline_name}Transform"
transformer = Transformer(
model_name=xgb_model.name,
strategy="SingleRecord",
instance_type="ml.m5.xlarge",
instance_count=1,
output_path=transform_output,
sagemaker_session=pipeline_session,
)

transform_input = pipeline_session.upload_data(
path=os.path.join(DATA_DIR, "xgboost_abalone", "abalone"),
key_prefix="integ-test-data/xgboost_abalone/abalone",
)

execution = transformer.transform_with_monitoring(
monitoring_config=data_bias_check_config,
monitoring_resource_config=check_job_config,
data=transform_input,
content_type="text/libsvm",
supplied_baseline_constraints=data_bias_supplied_baseline_constraints,
role=role,
fail_on_violation=False,
)

execution_steps = execution.list_steps()
assert len(execution_steps) == 2

for execution_step in execution_steps:
assert execution_step["StepStatus"] == "Succeeded"

xgb_model.delete_model()
61 changes: 60 additions & 1 deletion tests/unit/test_transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,6 +23,13 @@

from tests.integ import test_local_mode
from tests.unit import SAGEMAKER_CONFIG_TRANSFORM_JOB
from sagemaker.model_monitor import DatasetFormat
from sagemaker.workflow.quality_check_step import (
ModelQualityCheckConfig,
)
from sagemaker.workflow.check_job_config import CheckJobConfig

_CHECK_JOB_PREFIX = "CheckJobPrefix"

ROLE = "DummyRole"
REGION = "us-west-2"
Expand All@@ -49,6 +56,16 @@
"base_transform_job_name": JOB_NAME,
}

PROCESS_REQUEST_ARGS = {
"inputs": "processing_inputs",
"output_config": "output_config",
"job_name": "job_name",
"resources": "resource_config",
"stopping_condition": {"MaxRuntimeInSeconds": 3600},
"app_specification": "app_specification",
"experiment_config": {"ExperimentName": "AnExperiment"},
}

MODEL_DESC_PRIMARY_CONTAINER = {"PrimaryContainer": {"Image": IMAGE_URI}}

MODEL_DESC_CONTAINERS_ONLY = {"Containers": [{"Image": IMAGE_URI}]}
Expand All@@ -72,7 +89,7 @@ def mock_create_tar_file():

@pytest.fixture()
def sagemaker_session():
boto_mock = Mock(name="boto_session")
boto_mock = Mock(name="boto_session", region_name=REGION)
session = Mock(
name="sagemaker_session",
boto_session=boto_mock,
Expand DownExpand Up@@ -764,6 +781,48 @@ def test_stop_transform_job(sagemaker_session, transformer):
sagemaker_session.stop_transform_job.assert_called_once_with(name=JOB_NAME)


@patch("sagemaker.transformer.Transformer._retrieve_image_uri", return_value=IMAGE_URI)
@patch("sagemaker.workflow.pipeline.Pipeline.upsert", return_value={})
@patch("sagemaker.workflow.pipeline.Pipeline.start", return_value=Mock())
def test_transform_with_monitoring_create_and_starts_pipeline(
pipeline_start, upsert, image_uri, sagemaker_session, transformer
):

config = CheckJobConfig(
role=ROLE,
instance_count=1,
instance_type="ml.m5.xlarge",
volume_size_in_gb=60,
max_runtime_in_seconds=1800,
sagemaker_session=sagemaker_session,
base_job_name=_CHECK_JOB_PREFIX,
)

quality_check_config = ModelQualityCheckConfig(
baseline_dataset="s3://baseline_dataset_s3_url",
dataset_format=DatasetFormat.csv(header=True),
problem_type="BinaryClassification",
inference_attribute="quality_cfg_attr_value",
probability_attribute="quality_cfg_attr_value",
ground_truth_attribute="quality_cfg_attr_value",
probability_threshold_attribute="quality_cfg_attr_value",
post_analytics_processor_script="s3://my_bucket/data_quality/postprocessor.py",
output_s3_uri="s3://output_s3_uri",
)

transformer.transform_with_monitoring(
monitoring_config=quality_check_config,
monitoring_resource_config=config,
data=DATA,
content_type="text/libsvm",
supplied_baseline_constraints="supplied_baseline_constraints",
role=ROLE,
)

upsert.assert_called_once()
pipeline_start.assert_called_once()


def test_stop_transform_job_no_transform_job(transformer):
with pytest.raises(ValueError) as e:
transformer.stop_transform_job()
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
7 changes: 4 additions & 3 deletions src/sagemaker/transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -337,6 +337,7 @@ def transform_with_monitoring(
wait: bool = True,
pipeline_name: str = None,
role: str = None,
fail_on_violation: bool = True,
):
"""Runs a transform job with monitoring job.

Expand All@@ -352,7 +353,6 @@ def transform_with_monitoring(
]): the monitoring configuration used for run model monitoring.
monitoring_resource_config (`sagemaker.workflow.check_job_config.CheckJobConfig`):
the check job (processing job) cluster resource configuration.
transform_step_args (_JobStepArguments): the transform step transform arguments.
data (str): Input data location in S3 for the transform job
data_type (str): What the S3 location defines (default: 'S3Prefix').
Valid values:
Expand DownExpand Up@@ -400,8 +400,6 @@ def transform_with_monitoring(
monitor_before_transform (bgool): If to run data quality
or model explainability monitoring type,
a true value of this flag indicates running the check step before the transform job.
fail_on_violation (Union[bool, PipelineVariable]): A opt-out flag to not to fail the
check step when a violation is detected.
supplied_baseline_statistics (Union[str, PipelineVariable]): The S3 path
to the supplied statistics object representing the statistics JSON file
which will be used for drift to check (default: None).
Expand All@@ -411,6 +409,8 @@ def transform_with_monitoring(
wait (bool): To determine if needed to wait for the pipeline execution to complete
pipeline_name (str): The name of the Pipeline for the monitoring and transfrom step
role (str): Execution role
fail_on_violation (Union[bool, PipelineVariable]): A opt-out flag to not to fail the
check step when a violation is detected.
"""

transformer = self
Expand DownExpand Up@@ -454,6 +454,7 @@ def transform_with_monitoring(
monitor_before_transform=monitor_before_transform,
supplied_baseline_constraints=supplied_baseline_constraints,
supplied_baseline_statistics=supplied_baseline_statistics,
fail_on_violation=fail_on_violation,
Comment thread
keshav-chandak marked this conversation as resolved.
)

pipeline_name = (
Expand Down
64 changes: 64 additions & 0 deletions tests/integ/test_transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -709,3 +709,67 @@ def test_transformer_and_monitoring_job(
assert execution_step["StepStatus"] == "Succeeded"

xgb_model.delete_model()


def test_transformer_and_monitoring_job_to_pass_with_no_failure_in_violation(
pipeline_session,
sagemaker_session,
role,
pipeline_name,
check_job_config,
data_bias_check_config,
):
xgb_model_data_s3 = pipeline_session.upload_data(
path=os.path.join(os.path.join(DATA_DIR, "xgboost_abalone"), "xgb_model.tar.gz"),
key_prefix="integ-test-data/xgboost/model",
)
data_bias_supplied_baseline_constraints = Constraints.from_file_path(
constraints_file_path=os.path.join(
DATA_DIR, "pipeline/clarify_check_step/data_bias/bad_cases/analysis.json"
),
sagemaker_session=sagemaker_session,
).file_s3_uri

xgb_model = XGBoostModel(
model_data=xgb_model_data_s3,
framework_version="1.3-1",
role=role,
sagemaker_session=sagemaker_session,
entry_point=os.path.join(os.path.join(DATA_DIR, "xgboost_abalone"), "inference.py"),
enable_network_isolation=True,
)

xgb_model.deploy(_INSTANCE_COUNT, _INSTANCE_TYPE)

transform_output = f"s3://{sagemaker_session.default_bucket()}/{pipeline_name}Transform"
transformer = Transformer(
model_name=xgb_model.name,
strategy="SingleRecord",
instance_type="ml.m5.xlarge",
instance_count=1,
output_path=transform_output,
sagemaker_session=pipeline_session,
)

transform_input = pipeline_session.upload_data(
path=os.path.join(DATA_DIR, "xgboost_abalone", "abalone"),
key_prefix="integ-test-data/xgboost_abalone/abalone",
)

execution = transformer.transform_with_monitoring(
monitoring_config=data_bias_check_config,
monitoring_resource_config=check_job_config,
data=transform_input,
content_type="text/libsvm",
supplied_baseline_constraints=data_bias_supplied_baseline_constraints,
role=role,
fail_on_violation=False,
)

execution_steps = execution.list_steps()
assert len(execution_steps) == 2

for execution_step in execution_steps:
assert execution_step["StepStatus"] == "Succeeded"

xgb_model.delete_model()
61 changes: 60 additions & 1 deletion tests/unit/test_transformer.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,6 +23,13 @@

from tests.integ import test_local_mode
from tests.unit import SAGEMAKER_CONFIG_TRANSFORM_JOB
from sagemaker.model_monitor import DatasetFormat
from sagemaker.workflow.quality_check_step import (
ModelQualityCheckConfig,
)
from sagemaker.workflow.check_job_config import CheckJobConfig

_CHECK_JOB_PREFIX = "CheckJobPrefix"

ROLE = "DummyRole"
REGION = "us-west-2"
Expand All@@ -49,6 +56,16 @@
"base_transform_job_name": JOB_NAME,
}

PROCESS_REQUEST_ARGS = {
"inputs": "processing_inputs",
"output_config": "output_config",
"job_name": "job_name",
"resources": "resource_config",
"stopping_condition": {"MaxRuntimeInSeconds": 3600},
"app_specification": "app_specification",
"experiment_config": {"ExperimentName": "AnExperiment"},
}

MODEL_DESC_PRIMARY_CONTAINER = {"PrimaryContainer": {"Image": IMAGE_URI}}

MODEL_DESC_CONTAINERS_ONLY = {"Containers": [{"Image": IMAGE_URI}]}
Expand All@@ -72,7 +89,7 @@ def mock_create_tar_file():

@pytest.fixture()
def sagemaker_session():
boto_mock = Mock(name="boto_session")
boto_mock = Mock(name="boto_session", region_name=REGION)
session = Mock(
name="sagemaker_session",
boto_session=boto_mock,
Expand DownExpand Up@@ -764,6 +781,48 @@ def test_stop_transform_job(sagemaker_session, transformer):
sagemaker_session.stop_transform_job.assert_called_once_with(name=JOB_NAME)


@patch("sagemaker.transformer.Transformer._retrieve_image_uri", return_value=IMAGE_URI)
@patch("sagemaker.workflow.pipeline.Pipeline.upsert", return_value={})
@patch("sagemaker.workflow.pipeline.Pipeline.start", return_value=Mock())
def test_transform_with_monitoring_create_and_starts_pipeline(
pipeline_start, upsert, image_uri, sagemaker_session, transformer
):

config = CheckJobConfig(
role=ROLE,
instance_count=1,
instance_type="ml.m5.xlarge",
volume_size_in_gb=60,
max_runtime_in_seconds=1800,
sagemaker_session=sagemaker_session,
base_job_name=_CHECK_JOB_PREFIX,
)

quality_check_config = ModelQualityCheckConfig(
baseline_dataset="s3://baseline_dataset_s3_url",
dataset_format=DatasetFormat.csv(header=True),
problem_type="BinaryClassification",
inference_attribute="quality_cfg_attr_value",
probability_attribute="quality_cfg_attr_value",
ground_truth_attribute="quality_cfg_attr_value",
probability_threshold_attribute="quality_cfg_attr_value",
post_analytics_processor_script="s3://my_bucket/data_quality/postprocessor.py",
output_s3_uri="s3://output_s3_uri",
)

transformer.transform_with_monitoring(
monitoring_config=quality_check_config,
monitoring_resource_config=config,
data=DATA,
content_type="text/libsvm",
supplied_baseline_constraints="supplied_baseline_constraints",
role=ROLE,
)

upsert.assert_called_once()
pipeline_start.assert_called_once()


def test_stop_transform_job_no_transform_job(transformer):
with pytest.raises(ValueError) as e:
transformer.stop_transform_job()
Expand Down