* update: record tallied results to GCS bucket

* fix: test percent

* update: accumulator support

* fix: use args for where to store results

* fix: args for results file

* fix: gs bucket ptrfix

* fix: tuning args

* fix: ci/cd test on openin artifacts bucket

* fix: 2nd try at bucket issue

* fix: 2nd try at bucket issue

* debug: bucket issue

* debug: bucket issue

* debug: bucket issue

* debug: bucket issue

* debug: bucket issue

* debug: bucket issue

* debug: no entries written

* debug: no entries written

* debug: not accumulating

* debug: not accumulating

* debug: not accumulating

* debug: not accumulating

* Update requirements.txt

* debug: not accumulating

* debug: not accumulating

* debug: matching notebook name

* debug: pandas problem

* debug: import issues

* debug: load

* Update requirements.txt

* Update requirements.txt

* debug: read csv

* debug: read csv

* debug: read csv

* debug: indxer

* debug: indexer

* debug: accum

* debug: accum

* debug: accum

* debug: accum

* debug: accum

* debug: accum

* debug: accum

* debug: casting

* debug: casting

* debug: casting

* debug: casting

* debug: duration nit

* debug: duration nit

* fix: flaky

* fix: BUILD_ID

* fix: BUILD_ID

* fix: BUILD_ID

* fix: BUILD_ID

* fix: BUILD_ID

* fix: BUILD_ID

* fix: BUILD_ID

* fix: review

* fix: review

* fix: review

* fix: biz logic

* fix: biz logic

* fix: review

* fix: build_id required

* fix: use json format

* review: JSON simplifing

* review: JSON simplifing

* update: NotbookExecutionResult updates

* fix: revert to JSON

* update: use util to read from bucket

* fix: remove pandas inmport
This commit is contained in:
Andrew Ferlitsch
2023-05-05 20:03:49 +00:00
committed by GitHub
parent 2b4f7834b8
commit 6195e7bbf9
4 changed files with 123 additions and 12 deletions
+19 -4
View File
@@ -17,7 +17,7 @@
import argparse import argparse
import pathlib import pathlib
import random import os
import execute_changed_notebooks_helper import execute_changed_notebooks_helper
@@ -47,6 +47,12 @@ parser.add_argument(
required=False, required=False,
default=100, default=100,
) )
parser.add_argument(
"--build_id",
type=str,
help="The build id (which may be a Cloud Build job specific or user explicit.",
required=True
)
parser.add_argument( parser.add_argument(
"--base_branch", "--base_branch",
help="The base git branch to diff against to find changed files.", help="The base git branch to diff against to find changed files.",
@@ -129,26 +135,35 @@ changed_notebooks = execute_changed_notebooks_helper.get_changed_notebooks(
base_branch=args.base_branch, base_branch=args.base_branch,
) )
results_bucket = f"{args.artifacts_bucket}"
results_file = f"{args.build_id}.json"
if args.test_percent == 100: if args.test_percent == 100:
notebooks = changed_notebooks notebooks = changed_notebooks
accumulative_results = {}
else: else:
notebooks = [changed_notebook for changed_notebook in changed_notebooks if random.randint(1, 100) < args.test_percent] accumulative_results = execute_changed_notebooks_helper.load_results(results_bucket, results_file)
notebooks = [changed_notebook for changed_notebook in changed_notebooks if execute_changed_notebooks_helper.select_notebook(changed_notebook, accumulative_results, args.test_percent)]
if args.dry_run: if args.dry_run:
print("Dry run ...\n") print("Dry run ...\n")
for notebook in notebooks: for notebook in notebooks:
print(f"Would execute: {notebook.path}") print(f"Would execute: {notebook}")
else: else:
execute_changed_notebooks_helper.process_and_execute_notebooks( execute_changed_notebooks_helper.process_and_execute_notebooks(
notebooks=notebooks, notebooks=notebooks,
container_uri=args.container_uri, container_uri=args.container_uri,
staging_bucket=args.staging_bucket, staging_bucket=args.staging_bucket,
artifacts_bucket=args.artifacts_bucket, artifacts_bucket=args.artifacts_bucket,
results_file=results_file,
accumulative_results=accumulative_results,
should_parallelize=args.should_parallelize, should_parallelize=args.should_parallelize,
timeout=args.timeout, timeout=args.timeout,
variable_project_id=args.variable_project_id, variable_project_id=args.variable_project_id,
variable_region=args.variable_region, variable_region=args.variable_region,
variable_service_account=args.variable_service_account, variable_service_account=args.variable_service_account,
variable_vpc_network=args.variable_vpc_network, variable_vpc_network=args.variable_vpc_network,
private_pool_id=args.private_pool_id, private_pool_id=args.private_pool_id
) )
@@ -21,11 +21,15 @@ import json
import git import git
import operator import operator
import os import os
import io
import json
import pathlib import pathlib
import re import re
import subprocess import subprocess
import random
from google.cloud import storage
import utils import utils
from typing import List, Optional from typing import List, Optional, Dict, Any
from utils import util from utils import util
import execute_notebook_helper import execute_notebook_helper
@@ -65,7 +69,9 @@ def format_timedelta(delta: datetime.timedelta) -> str:
@dataclasses.dataclass @dataclasses.dataclass
class NotebookExecutionResult: class NotebookExecutionResult:
name: str name: str
path: str
duration: datetime.timedelta duration: datetime.timedelta
start_time: datetime.datetime
is_pass: bool is_pass: bool
log_url: str log_url: str
output_uri: str output_uri: str
@@ -80,6 +86,40 @@ class NotebookExecutionResult:
else: else:
return None return None
def load_results(results_bucket: str,
results_file: str) -> Dict[str,Any]:
'''
Load accumulated notebook test results
'''
print("Loading existing accumulative results ...")
accumulative_results = {}
try:
content = util.download_blob_into_memory(results_bucket, results_file, download_as_text=True)
accumulative_results = json.loads(content)
print(accumulative_results)
except Exception as e:
print(e)
# If there are no accumulative results, an empty dict is returned
return accumulative_results
def select_notebook(changed_notebook: str,
accumulative_results: Dict[str, Any],
test_percent: int) -> bool:
'''
Algorithm to randomly select a notebook, but weight the propbability of selected based on past failures
'''
if changed_notebook in accumulative_results:
pass_count = accumulative_results[changed_notebook]['passed']
fail_count = accumulative_results[changed_notebook]['failed']
else:
pass_count = 1
fail_count = 0
return (random.randint(1, 100) * (1 + (fail_count / (pass_count + fail_count))) < test_percent)
def _process_notebook( def _process_notebook(
notebook_path: str, notebook_path: str,
@@ -191,7 +231,9 @@ def process_and_execute_notebook(
result = NotebookExecutionResult( result = NotebookExecutionResult(
name=tag, name=tag,
path=notebook,
duration=datetime.timedelta(seconds=0), duration=datetime.timedelta(seconds=0),
start_time=datetime.datetime.now(),
is_pass=False, is_pass=False,
output_uri=notebook_output_uri, output_uri=notebook_output_uri,
log_url="", log_url="",
@@ -201,7 +243,6 @@ def process_and_execute_notebook(
) )
# TODO: Handle cases where multiple notebooks have the same name # TODO: Handle cases where multiple notebooks have the same name
time_start = datetime.datetime.now()
operation = None operation = None
try: try:
# Get the python version for running the notebook if specified # Get the python version for running the notebook if specified
@@ -247,9 +288,10 @@ def process_and_execute_notebook(
# Block and wait for the result # Block and wait for the result
operation_result = operation.result(timeout=timeout_in_seconds) operation_result = operation.result(timeout=timeout_in_seconds)
result.duration = datetime.datetime.now() - time_start result.duration = datetime.datetime.now() - result.start_time
result.is_pass = True result.is_pass = True
print(f"{notebook} PASSED in {format_timedelta(result.duration)}.") print(f"{notebook} PASSED in {format_timedelta(result.duration)}.")
except Exception as error: except Exception as error:
result.error_message = str(error) result.error_message = str(error)
@@ -268,7 +310,7 @@ def process_and_execute_notebook(
except Exception as error: except Exception as error:
result.error_message = str(error) result.error_message = str(error)
result.duration = datetime.datetime.now() - time_start result.duration = datetime.datetime.now() - result.start_time
result.is_pass = False result.is_pass = False
print( print(
@@ -336,12 +378,54 @@ def get_changed_notebooks(
return notebooks return notebooks
def _save_results(results: List[NotebookExecutionResult],
accumulative_results: Dict[str,Any],
artifacts_bucket: str,
results_file: str):
artifacts_bucket = artifacts_bucket.replace("gs://", "").split('/')[0]
print("Updating accumulative results ...")
for result in results:
if result.path in accumulative_results:
accumulative_results[result.path]['duration'] = result.duration.total_seconds()
accumulative_results[result.path]['start_time'] = str(result.start_time)
if result.is_pass:
accumulative_results[result.path]['passed'] += 1
else:
accumulative_results[result.path]['failed'] += 1
print(f"updating {result.path}")
else:
if result.is_pass:
pass_count = 1
fail_count = 0
else:
pass_count = 0
fail_count = 1
accumulative_results[result.path] = {
'duration': result.duration.total_seconds(),
'start_time': str(result.start_time),
'passed': pass_count,
'failed': fail_count
}
print(f"adding {result.path}")
print("Saving accumulative results ...")
content = json.dumps(accumulative_results)
client = storage.Client()
bucket = client.get_bucket(artifacts_bucket)
bucket.blob(str(results_file)).upload_from_string(content, 'text/json')
def process_and_execute_notebooks( def process_and_execute_notebooks(
notebooks: List[str], notebooks: List[str],
container_uri: str, container_uri: str,
staging_bucket: str, staging_bucket: str,
artifacts_bucket: str, artifacts_bucket: str,
results_file: str,
accumulative_results: List[NotebookExecutionResult],
should_parallelize: bool, should_parallelize: bool,
timeout: int, timeout: int,
variable_project_id: str, variable_project_id: str,
@@ -369,6 +453,10 @@ def process_and_execute_notebooks(
Required. The GCS staging bucket to write source code to. Required. The GCS staging bucket to write source code to.
artifacts_bucket (str): artifacts_bucket (str):
Required. The GCS staging bucket to write executed notebooks to. Required. The GCS staging bucket to write executed notebooks to.
results_file (str):
Required: The path to the artifacts bucket to save results
accumulative_results (List):
Required: The in-memory previous accumulative notebook CI/CD test results.
variable_project_id (str): variable_project_id (str):
Required. The value for PROJECT_ID to inject into notebooks. Required. The value for PROJECT_ID to inject into notebooks.
variable_region (str): variable_region (str):
@@ -471,7 +559,7 @@ def process_and_execute_notebooks(
print("=" * 100) print("=" * 100)
build_id = results_sorted[0].build_id build_id = results_sorted[0].build_id
logs_bucket_name = (results_sorted[0].logs_bucket).removeprefix("gs://") logs_bucket_name = (results_sorted[0].logs_bucket).replace("gs://", "")
log_file_name = f"log-{build_id}.txt" log_file_name = f"log-{build_id}.txt"
log_contents = util.download_blob_into_memory( log_contents = util.download_blob_into_memory(
@@ -489,6 +577,11 @@ def process_and_execute_notebooks(
else: else:
print(log_contents) print(log_contents)
_save_results(results_sorted,
accumulative_results,
artifacts_bucket,
results_file)
print("\n=== END RESULTS===\n") print("\n=== END RESULTS===\n")
total_notebook_duration = functools.reduce( total_notebook_duration = functools.reduce(
@@ -36,7 +36,7 @@ steps:
- -c - -c
- | - |
. workspace/env/bin/activate && . workspace/env/bin/activate &&
python3 .cloud-build/execute_changed_notebooks_cli.py --test_paths_file "${_TEST_PATHS_FILE}" --base_branch "${_FORCED_BASE_BRANCH}" --container_uri ${_PYTHON_IMAGE} --staging_bucket ${_GCS_STAGING_BUCKET} --artifacts_bucket ${_GCS_STAGING_BUCKET}/executed_notebooks/PR_${_PR_NUMBER}/BUILD_${BUILD_ID} --variable_project_id ${PROJECT_ID} --variable_region ${_GCP_REGION} --variable_service_account ${_GCP_SERVICE_ACCOUNT} --variable_vpc_network "${_GPC_VPC_NETWORK_NAME}" `if [ ! -z "${_PRIVATE_POOL_NAME}" ]; then echo "--private_pool_id ${_PRIVATE_POOL_NAME}"; fi` python3 .cloud-build/execute_changed_notebooks_cli.py --test_paths_file "${_TEST_PATHS_FILE}" --base_branch "${_FORCED_BASE_BRANCH}" --container_uri ${_PYTHON_IMAGE} --staging_bucket ${_GCS_STAGING_BUCKET} --artifacts_bucket ${_GCS_STAGING_BUCKET}/executed_notebooks/PR_${_PR_NUMBER}/BUILD_${BUILD_ID} --variable_project_id ${PROJECT_ID} --variable_region ${_GCP_REGION} --variable_service_account ${_GCP_SERVICE_ACCOUNT} --variable_vpc_network "${_GPC_VPC_NETWORK_NAME}" `if [ ! -z "${_PRIVATE_POOL_NAME}" ]; then echo "--private_pool_id ${_PRIVATE_POOL_NAME}"; fi` --build_id ${BUILD_ID}
env: env:
- 'IS_TESTING=1' - 'IS_TESTING=1'
timeout: 86400s timeout: 86400s
+5 -2
View File
@@ -3,12 +3,15 @@ numpy
jupyter jupyter
nbconvert nbconvert
papermill papermill
pandas
matplotlib matplotlib
tabulate tabulate
google-cloud-aiplatform google-cloud-aiplatform
google-cloud-storage google-cloud-storage
google-cloud-build google-cloud-build
google-cloud-storage
ratemate ratemate
GitPython GitPython
tqdm tqdm
fsspec
pandas