mirror of
https://github.com/encounter/adk-python.git
synced 2026-07-09 18:19:28 -07:00
feat: Add get_job_info tool to BigQuery toolset
This CL introduces a new tool, get_job_info, to the BigQuery toolset. This tool allows retrieving metadata about a BigQuery job, such as slot usage, job configuration, statistics, and job status. Closes #2928 Co-authored-by: Dongyu Jia <dongyuj@google.com> PiperOrigin-RevId: 825762399
This commit is contained in:
committed by
Copybara-Service
parent
72a8d8d85b
commit
64294572c1
@@ -21,6 +21,9 @@ distributed via the `google.adk.tools.bigquery` module. These tools include:
|
||||
|
||||
Fetches metadata about a BigQuery table.
|
||||
|
||||
5. `get_job_info`
|
||||
Fetches metadata about a BigQuery job.
|
||||
|
||||
5. `execute_sql`
|
||||
|
||||
Runs or dry-runs a SQL query in BigQuery.
|
||||
|
||||
@@ -80,6 +80,7 @@ class BigQueryToolset(BaseToolset):
|
||||
metadata_tool.get_table_info,
|
||||
metadata_tool.list_dataset_ids,
|
||||
metadata_tool.list_table_ids,
|
||||
metadata_tool.get_job_info,
|
||||
query_tool.get_execute_sql(self._tool_settings),
|
||||
query_tool.forecast,
|
||||
query_tool.analyze_contribution,
|
||||
|
||||
@@ -297,3 +297,297 @@ def get_table_info(
|
||||
"status": "ERROR",
|
||||
"error_details": str(ex),
|
||||
}
|
||||
|
||||
|
||||
def get_job_info(
|
||||
project_id: str,
|
||||
job_id: str,
|
||||
credentials: Credentials,
|
||||
settings: BigQueryToolConfig,
|
||||
) -> dict:
|
||||
"""Get metadata information about a BigQuery job. Including slot usage,
|
||||
job configuration, job statistics, job status, original query etc.
|
||||
|
||||
Args:
|
||||
project_id (str): The Google Cloud project id containing the job.
|
||||
job_id (str): The BigQuery job id.
|
||||
credentials (Credentials): The credentials to use for the request.
|
||||
settings (BigQueryToolConfig): The BigQuery tool settings.
|
||||
|
||||
Returns:
|
||||
dict: Dictionary representing the properties of the job.
|
||||
|
||||
Examples:
|
||||
>>> user may give job id in fomat of: project_id:region.job_id
|
||||
like bigquery-public-data:US.bquxjob_12345678_1234567890
|
||||
>>> get_job_info("bigquery-public-data", "bquxjob_12345678_1234567890")
|
||||
{
|
||||
"get_job_info_response": {
|
||||
"configuration": {
|
||||
"jobType": "QUERY",
|
||||
"query": {
|
||||
"destinationTable": {
|
||||
"datasetId": "_fd6de55d5d5c13fcfb0449cbf933bb695b2c3085",
|
||||
"projectId": "projectid",
|
||||
"tableId": "anonfbbe65d6_9782_469b_9f56_1392560314b2"
|
||||
},
|
||||
"priority": "INTERACTIVE",
|
||||
"query": "SELECT * FROM `projectid.dataset_id.table_id` WHERE TIMESTAMP_TRUNC(_PARTITIONTIME, DAY) = TIMESTAMP(\"2025-10-29\") LIMIT 1000",
|
||||
"useLegacySql": false,
|
||||
"writeDisposition": "WRITE_TRUNCATE"
|
||||
}
|
||||
},
|
||||
"etag": "EdeYv9sdcO7tD9HsffvcuQ==",
|
||||
"id": "projectid:US.job-id",
|
||||
"jobCreationReason": {
|
||||
"code": "REQUESTED"
|
||||
},
|
||||
"jobReference": {
|
||||
"jobId": "job-id",
|
||||
"location": "US",
|
||||
"projectId": "projectid"
|
||||
},
|
||||
"kind": "bigquery#job",
|
||||
"principal_subject": "user:abc@google.com",
|
||||
"selfLink": "https://bigquery.googleapis.com/bigquery/v2/projects/projectid/jobs/job-id?location=US",
|
||||
"statistics": {
|
||||
"creationTime": 1761760370152,
|
||||
"endTime": 1761760371250,
|
||||
"finalExecutionDurationMs": "489",
|
||||
"query": {
|
||||
"billingTier": 1,
|
||||
"cacheHit": false,
|
||||
"estimatedBytesProcessed": "5597805",
|
||||
"metadataCacheStatistics": {
|
||||
"tableMetadataCacheUsage": [
|
||||
{
|
||||
"explanation": "Table does not have CMETA.",
|
||||
"tableReference": {
|
||||
"datasetId": "datasetId",
|
||||
"projectId": "projectid",
|
||||
"tableId": "tableId"
|
||||
},
|
||||
"unusedReason": "OTHER_REASON"
|
||||
}
|
||||
]
|
||||
},
|
||||
"queryPlan": [
|
||||
{
|
||||
"completedParallelInputs": "3",
|
||||
"computeMode": "BIGQUERY",
|
||||
"computeMsAvg": "13",
|
||||
"computeMsMax": "15",
|
||||
"computeRatioAvg": 0.054852320675105488,
|
||||
"computeRatioMax": 0.063291139240506333,
|
||||
"endMs": "1761760370422",
|
||||
"id": "0",
|
||||
"name": "S00: Input",
|
||||
"parallelInputs": "8",
|
||||
"readMsAvg": "18",
|
||||
"readMsMax": "21",
|
||||
"readRatioAvg": 0.0759493670886076,
|
||||
"readRatioMax": 0.088607594936708861,
|
||||
"recordsRead": "1690",
|
||||
"recordsWritten": "1690",
|
||||
"shuffleOutputBytes": "1031149",
|
||||
"shuffleOutputBytesSpilled": "0",
|
||||
"slotMs": "157",
|
||||
"startMs": "1761760370388",
|
||||
"status": "COMPLETE",
|
||||
"steps": [
|
||||
{
|
||||
"kind": "READ",
|
||||
"substeps": [
|
||||
"$2:extendedFields.$is_not_null, $3:extendedFields.traceId, $4:span.$is_not_null, $5:span.spanKind, $6:span.endTime, $7:span.startTime, $8:span.parentSpanId, $9:span.spanId, $10:span.name, $11:span.childSpanCount.$is_not_null, $12:span.childSpanCount.value, $13:span.sameProcessAsParentSpan.$is_not_null, $14:span.sameProcessAsParentSpan.value, $15:span.status.$is_not_null, $16:span.status.message, $17:span.status.code",
|
||||
"FROM projectid.dataset_id.table_id",
|
||||
"WHERE equal(timestamp_trunc($1, 3), 1761696000.000000000)"
|
||||
]
|
||||
},
|
||||
{
|
||||
"kind": "LIMIT",
|
||||
"substeps": [
|
||||
"1000"
|
||||
]
|
||||
},
|
||||
{
|
||||
"kind": "WRITE",
|
||||
"substeps": [
|
||||
"$2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17",
|
||||
"TO __stage00_output"
|
||||
]
|
||||
}
|
||||
],
|
||||
"waitMsAvg": "1",
|
||||
"waitMsMax": "1",
|
||||
"waitRatioAvg": 0.0042194092827004216,
|
||||
"waitRatioMax": 0.0042194092827004216,
|
||||
"writeMsAvg": "2",
|
||||
"writeMsMax": "2",
|
||||
"writeRatioAvg": 0.0084388185654008432,
|
||||
"writeRatioMax": 0.0084388185654008432
|
||||
},
|
||||
{
|
||||
"completedParallelInputs": "1",
|
||||
"computeMode": "BIGQUERY",
|
||||
"computeMsAvg": "22",
|
||||
"computeMsMax": "22",
|
||||
"computeRatioAvg": 0.092827004219409287,
|
||||
"computeRatioMax": 0.092827004219409287,
|
||||
"endMs": "1761760370428",
|
||||
"id": "1",
|
||||
"inputStages": [
|
||||
"0"
|
||||
],
|
||||
"name": "S01: Compute+",
|
||||
"parallelInputs": "1",
|
||||
"readMsAvg": "0",
|
||||
"readMsMax": "0",
|
||||
"readRatioAvg": 0,
|
||||
"readRatioMax": 0,
|
||||
"recordsRead": "1001",
|
||||
"recordsWritten": "1000",
|
||||
"shuffleOutputBytes": "800157",
|
||||
"shuffleOutputBytesSpilled": "0",
|
||||
"slotMs": "29",
|
||||
"startMs": "1761760370398",
|
||||
"status": "COMPLETE",
|
||||
"steps": [
|
||||
{
|
||||
"kind": "READ",
|
||||
"substeps": [
|
||||
"$2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17",
|
||||
"FROM __stage00_output"
|
||||
]
|
||||
},
|
||||
{
|
||||
"kind": "COMPUTE",
|
||||
"substeps": [
|
||||
"$130 := MAKE_STRUCT($3, $2)",
|
||||
"$131 := MAKE_STRUCT($10, $9, $8, MAKE_STRUCT($29, $28, $27), $7, $6, MAKE_STRUCT(...), MAKE_STRUCT(...), MAKE_STRUCT(...), ...)"
|
||||
]
|
||||
},
|
||||
{
|
||||
"kind": "LIMIT",
|
||||
"substeps": [
|
||||
"1000"
|
||||
]
|
||||
},
|
||||
{
|
||||
"kind": "WRITE",
|
||||
"substeps": [
|
||||
"$130, $131",
|
||||
"TO __stage01_output"
|
||||
]
|
||||
}
|
||||
],
|
||||
"waitMsAvg": "7",
|
||||
"waitMsMax": "7",
|
||||
"waitRatioAvg": 0.029535864978902954,
|
||||
"waitRatioMax": 0.029535864978902954,
|
||||
"writeMsAvg": "4",
|
||||
"writeMsMax": "4",
|
||||
"writeRatioAvg": 0.016877637130801686,
|
||||
"writeRatioMax": 0.016877637130801686
|
||||
},
|
||||
{
|
||||
"completedParallelInputs": "1",
|
||||
"computeMode": "BIGQUERY",
|
||||
"computeMsAvg": "33",
|
||||
"computeMsMax": "33",
|
||||
"computeRatioAvg": 0.13924050632911392,
|
||||
"computeRatioMax": 0.13924050632911392,
|
||||
"endMs": "1761760370745",
|
||||
"id": "2",
|
||||
"inputStages": [
|
||||
"1"
|
||||
],
|
||||
"name": "S02: Output",
|
||||
"parallelInputs": "1",
|
||||
"readMsAvg": "0",
|
||||
"readMsMax": "0",
|
||||
"readRatioAvg": 0,
|
||||
"readRatioMax": 0,
|
||||
"recordsRead": "1000",
|
||||
"recordsWritten": "1000",
|
||||
"shuffleOutputBytes": "459829",
|
||||
"shuffleOutputBytesSpilled": "0",
|
||||
"slotMs": "106",
|
||||
"startMs": "1761760370667",
|
||||
"status": "COMPLETE",
|
||||
"steps": [
|
||||
{
|
||||
"kind": "READ",
|
||||
"substeps": [
|
||||
"$130, $131",
|
||||
"FROM __stage01_output"
|
||||
]
|
||||
},
|
||||
{
|
||||
"kind": "WRITE",
|
||||
"substeps": [
|
||||
"$130, $131",
|
||||
"TO __stage02_output"
|
||||
]
|
||||
}
|
||||
],
|
||||
"waitMsAvg": "237",
|
||||
"waitMsMax": "237",
|
||||
"waitRatioAvg": 1,
|
||||
"waitRatioMax": 1,
|
||||
"writeMsAvg": "55",
|
||||
"writeMsMax": "55",
|
||||
"writeRatioAvg": 0.2320675105485232,
|
||||
"writeRatioMax": 0.2320675105485232
|
||||
}
|
||||
],
|
||||
"referencedTables": [
|
||||
{
|
||||
"datasetId": "dataset_id",
|
||||
"projectId": "projectid",
|
||||
"tableId": "table_id"
|
||||
}
|
||||
],
|
||||
"statementType": "SELECT",
|
||||
"timeline": [
|
||||
{
|
||||
"completedUnits": "5",
|
||||
"elapsedMs": "492",
|
||||
"estimatedRunnableUnits": "0",
|
||||
"pendingUnits": "5",
|
||||
"totalSlotMs": "293"
|
||||
}
|
||||
],
|
||||
"totalBytesBilled": "10485760",
|
||||
"totalBytesProcessed": "5597805",
|
||||
"totalPartitionsProcessed": "2",
|
||||
"totalSlotMs": "293",
|
||||
"transferredBytes": "0"
|
||||
},
|
||||
"startTime": 1761760370268,
|
||||
"totalBytesProcessed": "5597805",
|
||||
"totalSlotMs": "293"
|
||||
},
|
||||
"status": {
|
||||
"state": "DONE"
|
||||
},
|
||||
"user_email": "abc@google.com"
|
||||
}
|
||||
}
|
||||
"""
|
||||
try:
|
||||
bq_client = client.get_bigquery_client(
|
||||
project=project_id,
|
||||
credentials=credentials,
|
||||
location=settings.location,
|
||||
user_agent=settings.application_name,
|
||||
)
|
||||
job = bq_client.get_job(job_id)
|
||||
# We need to use _properties to get the job info because it contains all
|
||||
# the job info.
|
||||
# pylint: disable=protected-access
|
||||
return job._properties
|
||||
except Exception as ex:
|
||||
return {
|
||||
"status": "ERROR",
|
||||
"error_details": str(ex),
|
||||
}
|
||||
|
||||
@@ -136,6 +136,36 @@ def test_get_table_info_no_default_auth(mock_default_auth, mock_get_table):
|
||||
mock_default_auth.assert_not_called()
|
||||
|
||||
|
||||
@mock.patch.dict(os.environ, {}, clear=True)
|
||||
@mock.patch("google.cloud.bigquery.Client.get_job", autospec=True)
|
||||
@mock.patch("google.auth.default", autospec=True)
|
||||
def test_get_job_info_no_default_auth(mock_default_auth, mock_get_job):
|
||||
"""Test get_job_info tool invocation involves no default auth."""
|
||||
mock_credentials = mock.create_autospec(Credentials, instance=True)
|
||||
tool_settings = BigQueryToolConfig()
|
||||
|
||||
# Simulate the behavior of default auth - on purpose throw exception when
|
||||
# the default auth is called
|
||||
mock_default_auth.side_effect = DefaultCredentialsError(
|
||||
"Your default credentials were not found"
|
||||
)
|
||||
|
||||
mock_get_job.return_value = mock.create_autospec(
|
||||
bigquery.QueryJob, instance=True
|
||||
)
|
||||
result = metadata_tool.get_job_info(
|
||||
"my_project_id",
|
||||
"my_job_id",
|
||||
mock_credentials,
|
||||
tool_settings,
|
||||
)
|
||||
assert result != {
|
||||
"status": "ERROR",
|
||||
"error_details": "Your default credentials were not found",
|
||||
}
|
||||
mock_default_auth.assert_not_called()
|
||||
|
||||
|
||||
@mock.patch(
|
||||
"google.adk.tools.bigquery.client.get_bigquery_client", autospec=True
|
||||
)
|
||||
|
||||
@@ -41,7 +41,7 @@ async def test_bigquery_toolset_tools_default():
|
||||
tools = await toolset.get_tools()
|
||||
assert tools is not None
|
||||
|
||||
assert len(tools) == 9
|
||||
assert len(tools) == 10
|
||||
assert all([isinstance(tool, GoogleTool) for tool in tools])
|
||||
|
||||
expected_tool_names = set([
|
||||
@@ -49,6 +49,7 @@ async def test_bigquery_toolset_tools_default():
|
||||
"get_dataset_info",
|
||||
"list_table_ids",
|
||||
"get_table_info",
|
||||
"get_job_info",
|
||||
"execute_sql",
|
||||
"ask_data_insights",
|
||||
"forecast",
|
||||
|
||||
Reference in New Issue
Block a user