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
222 changes: 222 additions & 0 deletions examples/job_results_paging.ipynb
Original file line number Diff line number Diff line change
@@ -0,0 +1,222 @@
{
"cells": [
{
"cell_type": "markdown",
"id": "088960ea",
"metadata": {},
"source": [
"## Harmony Py Library\n",
"### Paging Through Job Results\n",
"\n",
"`result_urls()` accepts optional `offset` and `limit` arguments so you can page through a\n",
"job's output links instead of pulling the complete list at once. It also accepts an\n",
"`allow_incomplete` flag to return whatever links are available so far, without waiting for\n",
"the job to reach a terminal state.\n",
"\n",
"This example uses `max_results=10` on the request so the job only produces a handful of\n",
"granules, keeping the paging examples below small and easy to read."
]
},
{
"cell_type": "code",
"execution_count": null,
"id": "72b6ffa9",
"metadata": {},
"outputs": [],
"source": [
"import sys\n",
"import helper\n",
"helper.install_project_and_dependencies('..')\n",
"\n",
"from harmony import BBox, Client, Collection, Request, Environment"
]
},
{
"cell_type": "markdown",
"id": "6a75edc9-df71-4791-80f3-d35e19fa1a0c",
"metadata": {},
"source": [
"#### Submit a request for processing and return the job_id"
]
},
{
"cell_type": "code",
"execution_count": null,
"id": "487dfa1d",
"metadata": {},
"outputs": [],
"source": [
"harmony_client = Client(env=Environment.UAT) # assumes .netrc usage\n",
"\n",
"collection = Collection(id='C1234088182-EEDTEST')\n",
"request = Request(\n",
" collection=collection,\n",
" spatial=BBox(-165, 52, -140, 77),\n",
" max_results=10, # keep the job small so paging is easy to see\n",
")\n",
"\n",
"job_id = harmony_client.submit(request)\n",
"job_id"
]
},
{
"cell_type": "markdown",
"id": "7d97b52e-ed9d-4948-aecf-879ef017fd46",
"metadata": {},
"source": [
"#### Default behavior: no offset/limit means every data URL is returned. Blocking - waits for the job to reach a terminal state."
]
},
{
"cell_type": "code",
"execution_count": null,
"id": "66d3573b",
"metadata": {},
"outputs": [],
"source": [
"all_urls = list(harmony_client.result_urls(job_id))\n",
"len(all_urls), all_urls"
]
},
{
"cell_type": "markdown",
"id": "b41726a4-8537-4eed-9681-86fe08d26cdb",
"metadata": {},
"source": [
"#### 'limit' caps how many URLs are returned"
]
},
{
"cell_type": "code",
"execution_count": null,
"id": "505af373",
"metadata": {},
"outputs": [],
"source": [
"first_page = list(harmony_client.result_urls(job_id, limit=3))\n",
"first_page"
]
},
{
"cell_type": "markdown",
"id": "f0703678-81e7-4c5f-a0d6-9b67f90e28d6",
"metadata": {},
"source": [
"#### 'page' skips past URLs already seen, so the next page picks up where the last one left off."
]
},
{
"cell_type": "code",
"execution_count": null,
"id": "1f704701",
"metadata": {},
"outputs": [],
"source": [
"second_page = list(harmony_client.result_urls(job_id, page=2, limit=3))\n",
"second_page"
]
},
{
"cell_type": "markdown",
"id": "085772ac-40d4-4fff-8733-15eac8dcaeb2",
"metadata": {},
"source": [
"#### A simple paging loop: keep requesting pages until an empty page comes back."
]
},
{
"cell_type": "code",
"execution_count": null,
"id": "7138253f",
"metadata": {},
"outputs": [],
"source": [
"page_size = 3\n",
"page = 1\n",
"while True:\n",
" results = list(harmony_client.result_urls(job_id, page=page, limit=page_size))\n",
" if not results:\n",
" break\n",
" print(f'page={page}: {results}')\n",
" page += 1"
]
},
{
"cell_type": "markdown",
"id": "26654e72-e39e-477c-b2ce-d4a0b22ff512",
"metadata": {},
"source": [
"#### Asking for a page/limit beyond the available results returns an empty list rather than raising an error."
]
},
{
"cell_type": "code",
"execution_count": null,
"id": "e04e5837",
"metadata": {},
"outputs": [],
"source": [
"list(harmony_client.result_urls(job_id, page=1000, limit=10))"
]
},
{
"cell_type": "markdown",
"id": "de249f1e-dfec-4780-aeff-e8da81863666",
"metadata": {},
"source": [
"#### 'allow_incomplete=True' returns whatever data URLs are currently available without waiting for the job to finish processing. This is useful for a still-running job when you just want to start downloading what's ready so far. Combine with offset/limit to page through the partial results too."
]
},
{
"cell_type": "code",
"execution_count": null,
"id": "dd018939",
"metadata": {},
"outputs": [],
"source": [
"partial_job_id = harmony_client.submit(request)\n",
"partial_job_id"
]
},
{
"cell_type": "markdown",
"id": "50f42c73-53a3-4b8c-b670-20a073a0529f",
"metadata": {},
"source": [
"#### Run this cell while the request is running to see the partial results returned without blocking"
]
},
{
"cell_type": "code",
"execution_count": null,
"id": "7568355a-8a72-4e31-8cfb-73ccef240f38",
"metadata": {},
"outputs": [],
"source": [
"partial_urls = list(harmony_client.result_urls(partial_job_id, allow_incomplete=True))\n",
"len(partial_urls), partial_urls"
]
}
],
"metadata": {
"kernelspec": {
"display_name": "Python 3 (ipykernel)",
"language": "python",
"name": "python3"
},
"language_info": {
"codemirror_mode": {
"name": "ipython",
"version": 3
},
"file_extension": ".py",
"mimetype": "text/x-python",
"name": "python",
"nbconvert_exporter": "python",
"pygments_lexer": "ipython3",
"version": "3.12.13"
}
},
"nbformat": 4,
"nbformat_minor": 5
}
89 changes: 78 additions & 11 deletions harmony/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -325,9 +325,29 @@ def _submit_url(self, request: BaseRequest) -> str:
else:
return request_url_map[type(request)](self, request)

def _status_url(self, job_id: str, link_type: LinkType = LinkType.https) -> str:
"""Constructs the URL for the Job that is used to get its status."""
return f'{self.config.root_url}/jobs/{job_id}?linktype={link_type.value}'
def _status_url(
self,
job_id: str,
link_type: LinkType = LinkType.https,
page: int = None,
limit: int = None,
) -> str:
"""Constructs the URL for the Job that is used to get its status.

Args:
job_id: UUID string for the job to be fetched
link_type: The type of link to output, s3:// or https://
page: The 1-indexed page of the job's output links to fetch, matching Harmony's
own ``page`` query parameter.
limit: The number of output links per page, matching Harmony's own ``limit``
query parameter.
"""
url = f'{self.config.root_url}/jobs/{job_id}?linktype={link_type.value}'
if page is not None:
url += f'&page={page}'
if limit is not None:
url += f'&limit={limit}'
return url

def _job_status_batch_url(self) -> str:
"""Constructs the URL for the lightweight batch job status endpoint."""
Expand Down Expand Up @@ -1052,44 +1072,91 @@ def _get_json(self, url: str) -> str:
return response.json()

def _result_pages(
self, job_id: str, show_progress: bool = False, link_type: LinkType = LinkType.https
self,
job_id: str,
show_progress: bool = False,
link_type: LinkType = LinkType.https,
allow_incomplete: bool = False,
page: int = None,
limit: int = None,
) -> Generator[object, None, None]:
"""Yields each page of results for the provided job ID

Args:
job_id: UUID string for the job to be fetched
show_progress: Whether a progress bar should show via stdout.
link_type: The type of link to output, s3:// or https://
allow_incomplete: If True, does not wait for the job to reach a terminal state,
and instead returns whatever result pages are currently available.
page: The 1-indexed page of output links to fetch, matching Harmony's own
``page`` query parameter. If provided, only this single page is yielded
rather than following every ``next`` link.
limit: The number of output links per page, matching Harmony's own ``limit``
query parameter.

Returns:
A generator for each page of results, loaded on demand
"""
self.wait_for_processing(job_id, show_progress)
next_url = self._status_url(job_id, link_type)
if not allow_incomplete:
self.wait_for_processing(job_id, show_progress)
next_url = self._status_url(job_id, link_type, page=page, limit=limit)
while next_url:
response = self._get_json(next_url)
yield response
if page is not None or limit is not None:
return
links = response.get('links', [])
next_url = next((x['href'] for x in links if x['rel'] == 'next'), None)

def result_urls(
self, job_id: str, show_progress: bool = False, link_type: LinkType = LinkType.https
self,
job_id: str,
show_progress: bool = False,
link_type: LinkType = LinkType.https,
page: int = None,
limit: int = None,
allow_incomplete: bool = False,
) -> Generator[str, None, None]:
"""Retrieve the data URLs for a job.

The URLs include links to all of the jobs data output. Blocks until the Harmony job is
done processing.
done processing, unless ``allow_incomplete`` is True.

Args:
job_id: UUID string for the job you wish to interrogate.
show_progress: Whether a progress bar should show via stdout.
link_type: The type of link to output, s3:// or https://
page: The 1-indexed page of output links to fetch, matching Harmony's own
``page`` query parameter. Must be an integer greater than 0. Defaults to
returning every page (i.e. all URLs), unless ``limit`` is provided, in which
case it defaults to page 1. If ``page``/``limit`` go beyond the available
results, an empty list is returned rather than an error.
limit: The number of output links per page, matching Harmony's own ``limit``
query parameter. Must be an integer between 1 and 2000, inclusive.
allow_incomplete: If True, returns the URLs available so far without waiting for the
job to reach a terminal state. Defaults to False, which only returns URLs once
the job has reached a terminal state.

Returns:
The job's complete list of data URLs.
The job's (optionally paged) list of data URLs.

Raises:
ValueError: If ``page`` is not an integer greater than 0, or ``limit`` is not an
integer between 1 and 2000 inclusive.
"""
for page in self._result_pages(job_id, show_progress, link_type):
for link in page.get('links', []):
if page is not None and (not isinstance(page, int) or isinstance(page, bool) or page < 1):
raise ValueError('page must be an integer greater than 0.')
if limit is not None and (
not isinstance(limit, int) or isinstance(limit, bool) or not 1 <= limit <= 2000
):
raise ValueError('limit must be an integer between 1 and 2000, inclusive.')
if page is None and limit is not None:
page = 1

for result_page in self._result_pages(
job_id, show_progress, link_type, allow_incomplete, page=page, limit=limit
):
for link in result_page.get('links', []):
if link['rel'] == 'data':
yield link['href']

Expand Down
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,7 @@ examples = [
# Version is derived from the git tag (e.g. v1.4.0 -> 1.4.0) at build time.

[tool.setuptools.packages.find]
exclude = ["contrib", "docs", "tests*"]
include = ["harmony", "harmony.*"]

[tool.ruff]
line-length = 99
Expand Down
Loading
Loading