diff --git a/Makefile b/Makefile index 8f3aceec..4deac042 100644 --- a/Makefile +++ b/Makefile @@ -58,7 +58,6 @@ geomad-test-ausp: --year 2000 \ --version 0-3-0-test \ --decimated \ - --no-single-region \ --bucket data.ldn.auspatious.com \ --overwrite; geomad-test-dep-staging: @@ -69,7 +68,6 @@ geomad-test-dep-staging: --version 0-3-0-test \ --collection-url-root="https://stac.staging.digitalearthpacific.io/collections" \ --decimated \ - --single-region \ --bucket dep-public-staging \ --overwrite; @@ -115,27 +113,28 @@ collection-geomad-test-dep-staging: -# #### Training Data -# training-data-generate: -# for site in $(PACIFIC_TRAINING_TILES); do \ -# tile_id=$$(echo $$site | cut -d: -f1); \ -# region=$$(echo $$site | cut -d: -f2); \ -# country_name=$$(echo $$site | cut -d: -f3 | tr '_' ' '); \ -# country_code=$$(echo $$site | cut -d: -f4); \ -# ldn training generate-training-data \ -# --tile-id $$tile_id \ -# --region $$region \ -# --country-name "$$country_name" \ -# --country-code "$$country_code"; \ -# done; - -# # poetry run ldn training generate-training-data \ -# # --tile-id 028_030 \ -# # --region pacific \ -# # --country-name "Papua New Guinea" \ -# # --country-code "PNG" \ -# # --geomad-version 0-2-1; +#### Training Data +# Geomad version: 0-2-1 in DEP staging, 0-3-0 in DEP public. +training-data-generate: + for site in $(PACIFIC_TRAINING_TILES) do \ + tile_id=$$(echo $$site | cut -d: -f1); \ + region=$$(echo $$site | cut -d: -f2); \ + country_name=$$(echo $$site | cut -d: -f3 | tr '_' ' '); \ + country_code=$$(echo $$site | cut -d: -f4); \ + ldn training generate-training-data \ + --tile-id $$tile_id \ + --region $$region \ + --country-name "$$country_name" \ + --country-code "$$country_code" \ + --geomad-version 0-2-1 \ + --geomad-bucket dep-public-staging \ + --output-bucket dep-public-staging \ + --single-region \ + --product-owner dep \ + --no-overwrite; \ + done; +#### Make the model using ldn-lulc/notebooks/1_Train_Model.ipynb @@ -143,7 +142,7 @@ collection-geomad-test-dep-staging: # # Predict LULC for the test tiles and one year (2025). -# # 1. Print tasks +# # Print tasks # print-tasks-lulc-2020: # ldn print-tasks \ # --years="2020" \ @@ -151,37 +150,39 @@ collection-geomad-test-dep-staging: # --dataset="lulc"; -# # 2. Classify -# predict-lulc-test-tiles-2020: -# for site in $(TEST_TILES); do \ -# tile_id=$${site%%:*}; \ -# region=$${site#*:}; region=$${region%%:*}; \ -# ldn lulc run \ -# --tile-id $$tile_id \ -# --year 2020 \ -# --version $(LULC_VERSION) \ -# --geomad-version $(GEOMAD_VERSION) \ -# --region $$region \ -# $(DECIMATED) \ -# --overwrite; \ -# done; - -# lulc-2-regions-decimated: -# for site in $(TEST_TILES_2_REGIONS); do \ -# tile_id=$${site%%:*}; \ -# region=$${site#*:}; region=$${region%%:*}; \ -# ldn lulc run \ -# --tile-id $$tile_id \ -# --year 2010 \ -# --version $(LULC_VERSION) \ -# --geomad-version $(GEOMAD_VERSION) \ -# --region $$region \ -# --decimated \ -# --overwrite; \ -# done; - - -# # 3. Update the STAC-Geoparquet index after all tiles/years have run. +# # Classify +lulc-predict-test: + ldn lulc run \ + --tile-id 028_030 \ + --year 2000 \ + --region pacific \ + --version 0-0-9 \ + --geomad-version 0-2-1 \ + --bucket dep-public-staging \ + --product-owner dep \ + --model-path="/Users/wj/Projects/ldn-lulc/ldn-lulc/ldn/models/0-0-9/pacific/2020/lulc_random_forest_model_pacific_2020.joblib" \ + --no-overwrite; + +# --model-path="https://dep-public-staging.s3.us-west-2.amazonaws.com/dep_ls_lulc/models/0-0-9/pacific/2020/lulc_random_forest_model_pacific_2020.joblib" \ + +lulc-predict-test-2: + for site in $(PACIFIC_TRAINING_TILES) do \ + tile_id=$$(echo $$site | cut -d: -f1); \ + region=$$(echo $$site | cut -d: -f2); \ + for year in 2000 2025; do \ + ldn lulc run \ + --tile-id $$tile_id \ + --year $$year \ + --region $$region \ + --version 0-0-9 \ + --geomad-version 0-2-1 \ + --bucket dep-public-staging \ + --product-owner dep \ + --model-path="/Users/wj/Projects/ldn-lulc/ldn-lulc/ldn/models/0-0-9/pacific/2020/lulc_random_forest_model_pacific_2020.joblib" \ + --no-overwrite; + done; + done; +# # Update the STAC-Geoparquet index after all tiles/years have run. # index-lulc: # ldn index-to-stac-geoparquet \ # --dataset "lulc" \ diff --git a/ldn/cli_geomad.py b/ldn/cli_geomad.py index 3c583429..d134080f 100644 --- a/ldn/cli_geomad.py +++ b/ldn/cli_geomad.py @@ -85,10 +85,6 @@ def run( year: Annotated[str, typer.Option()], version: Annotated[str, typer.Option()], region: Annotated[Literal["pacific", "non-pacific"], typer.Option()], - single_region: Annotated[ - bool, - typer.Option(help="Whether to use the single region prefix for the collection_url_root (e.g. 'dep_ls_geomad')"), - ], bucket: Annotated[str | None, typer.Option(help="S3 bucket for data.")] = None, product_owner: Annotated[str | None, typer.Option(help="Override the region-derived owner prefix.")] = None, overwrite: Annotated[bool, typer.Option()] = False, @@ -112,7 +108,7 @@ def run( str | None, typer.Option( help="Override the default collection URL root" - " e.g for a STAC API like 'https://stac.digitalearthpacific.org/collections/dep_ls_geomad'" + " e.g for a STAC API like 'https://stac.digitalearthpacific.org/collections'" ), ] = None, sensor: Annotated[str, typer.Option(help="Sensor name, e.g. 'ls'.")] = SENSOR, @@ -221,6 +217,7 @@ def run( overwrite, collection_url_root=collection_url_root, s3_client=s3_client, + sensor=sensor, ) if components is None: return # Skip due to no overwrite. diff --git a/ldn/cli_lulc.py b/ldn/cli_lulc.py index 384da3ed..bbff6130 100644 --- a/ldn/cli_lulc.py +++ b/ldn/cli_lulc.py @@ -7,10 +7,9 @@ from ldn.utils import ( GEOMAD_VERSION, LULC_VERSION, - MODEL_VERSION, + SENSOR, LdnError, get_env_var, - owner_for_region, ) classify_app = typer.Typer() @@ -36,7 +35,7 @@ def run( product_owner: str | None = typer.Option(None, help="Override the region-derived owner prefix."), model_path: str = typer.Option( # TODO: defaults to pacific. Later have per region/time period models. - f"https://s3.us-west-2.amazonaws.com/data.ldn.auspatious.com/models/{MODEL_VERSION}/pacific/2020/lulc_random_forest_model_pacific_2020.joblib", + "https://dep-public-staging.s3.us-west-2.amazonaws.com/dep_ls_lulc/models/0-0-9/pacific/2020/lulc_random_forest_model_pacific_2020.joblib", help="Model to use for LULC classification.", ), decimated: bool = typer.Option( @@ -66,18 +65,19 @@ def run( help="Chunk size in pixels for x and y dimensions. Larger chunk sizes may be faster but use more memory." ), ] = 1024, - single_region: bool = typer.Option( - ..., - help="Whether to use the single region prefix (e.g. 'dep_ls_geomad') " - "or the generic prefix (e.g. 'ls_geomad') when accessing GeoMAD data.", - ), + sensor: str = typer.Option(SENSOR, help="Sensor to use for LULC classification. Defaults to 'ls'."), + collection_url_root: Annotated[ + str | None, + typer.Option( + help="Override the default collection URL root" + " e.g for a STAC API like 'https://stac.digitalearthpacific.org/collections'" + ), + ] = None, ) -> None: if int(year) < 2000 or int(year) > 2025: raise LdnError("Year must be between 2000 and 2025.") bucket = bucket or get_env_var("BUCKET") # Default - owner = owner_for_region(region, product_owner) - # TODO: Use build_prefix() here? run_classify_task( tile_id, @@ -86,7 +86,7 @@ def run( geomad_version=geomad_version, region=region, bucket=bucket, - owner=owner, + product_owner=product_owner, model_path=model_path, xy_chunk_size=xy_chunk_size, decimated=decimated, @@ -97,5 +97,6 @@ def run( memory_limit=memory_limit, n_workers=n_workers, threads_per_worker=threads_per_worker, - single_region=single_region, + sensor=sensor, + collection_url_root=collection_url_root, ) diff --git a/ldn/lulc.py b/ldn/lulc.py index 227ce1d6..24aa2dba 100644 --- a/ldn/lulc.py +++ b/ldn/lulc.py @@ -22,7 +22,7 @@ from sklearn.ensemble import RandomForestClassifier from typing_extensions import Annotated -from ldn.aws import configure_s3_access_profile +from ldn.aws import configure_s3_access_profile, s3_client from ldn.geomad import AwsStacTask as Task from ldn.grids import get_gridspec from ldn.raster import ( @@ -33,6 +33,7 @@ scale_offset_landsat, ) from ldn.utils import ( + GEOMAD_DATASET_ID, GEOMAD_VERSION, LULC_DATASET_ID, LULC_VERSION, @@ -40,8 +41,11 @@ WGS84, LdnError, get_analysis_epsg, + get_public_url_base, + get_stac_geoparquet_key, get_stac_geoparquet_url, is_bucket_source_coop, + owner_for_region, parse_tile_id, ) @@ -206,7 +210,13 @@ def do_prediction( nodata_mask = stacked.isnull().any(dim="variable") # Build observation table: fill NaN with nodata_value (masked pixels are excluded below). - obs = stacked.squeeze().fillna(nodata_value).transpose().to_dataframe() + obs = ( + stacked.squeeze() + .fillna(nodata_value) + .transpose() + .to_dataset(dim="variable") # pivot "variable" into named columns + .to_dataframe() + ) # Validate that all model features are present before reindexing. missing = set(model.feature_names_in_) - set(obs.columns) @@ -382,7 +392,7 @@ def run_classify_task( geomad_version: Annotated[str, typer.Option()], region: Literal["pacific", "non-pacific"], bucket: str, - owner: str, + product_owner: str | None, model_path: str, xy_chunk_size: int, decimated: bool, @@ -393,7 +403,8 @@ def run_classify_task( memory_limit: str, n_workers: int, threads_per_worker: int, - single_region: bool, + sensor: str, + collection_url_root: str | None, ) -> None: """Run LULC prediction for a single tile and year, writing results to S3. @@ -407,7 +418,7 @@ def run_classify_task( geomad_version: Version of the GeoMAD data to use (e.g. "0-0-1"). region: Grid region, either "pacific" or "non-pacific". bucket: S3 bucket for output COGs, STAC metadata, and input GeoMAD source data. - owner: Output prefix for paths (e.g. "dep" or "ci" or owner override). + product_owner: Override the region-derived owner prefix. model_path: Path or URL to the trained joblib model. xy_chunk_size: Chunk size in pixels for lazy loading. decimated: If True, use 10x lower resolution (for testing). @@ -418,6 +429,9 @@ def run_classify_task( memory_limit: Per-worker Dask memory limit. n_workers: Number of Dask workers. threads_per_worker: Number of threads per Dask worker. + single_region: If True, use the single region prefix (e.g. 'dep_ls_geomad') for GeoMAD data. + sensor: Sensor to use for LULC classification e.g. 'ls'. + collection_url_root: Override the default collection URL root for STAC metadata. """ logger.info(f"Starting processing. Tile ID: {tile_id}, Year: {year}, Region: {region}, Version: {version}.") logger.info( @@ -433,11 +447,11 @@ def run_classify_task( logger.info( f"Overriding the latest LULC prediction version ({LULC_VERSION}) with the specified version ({version})." ) + owner = owner_for_region(region, product_owner) - # geomad_stac_geoparquet_key = get_stac_geoparquet_key( - # bucket, single_region, product_owner, sensor, "geomad", geomad_version - # ) - geomad_stac_geoparquet_key = "" # TODO: Fix this. + geomad_stac_geoparquet_key = get_stac_geoparquet_key( + bucket, product_owner, sensor, GEOMAD_DATASET_ID, geomad_version + ) geomad_stac_geoparquet_url = get_stac_geoparquet_url(bucket, geomad_stac_geoparquet_key) tile_id_tuple = parse_tile_id(tile_id) @@ -464,18 +478,30 @@ def run_classify_task( logger.info("Loading model") loaded_model = _load_joblib_model(model_path) + # TODO: Could this block be a function because it will be done in cli_collection.py? + _is_source_coop = is_bucket_source_coop(bucket) + public_url = get_public_url_base(bucket) + if _is_source_coop: + public_url = f"{public_url}/{SOURCE_COOP_PREFIX_LULC}" + collection_url_root = collection_url_root or f"{public_url}/collections" + + owner = owner_for_region(region, product_owner) + components = build_pipeline_components( tile_id_tuple, year, version, bucket, - owner, + owner, # This owner respects single_region because geomad always writes to an owner. LULC_DATASET_ID, - SOURCE_COOP_PREFIX_LULC if is_bucket_source_coop(bucket) else None, + SOURCE_COOP_PREFIX_LULC if _is_source_coop else None, overwrite, + collection_url_root=collection_url_root, + s3_client=s3_client, + sensor=sensor, ) if components is None: - return # Task exists and overwrite is False, so skipping processing. + return # Skip due to no overwrite. itempath, stac_creator, writer = components searcher = StacGeoparquetSearcher( diff --git a/ldn/raster.py b/ldn/raster.py index f9f4d60f..d3be0a27 100644 --- a/ldn/raster.py +++ b/ldn/raster.py @@ -19,7 +19,6 @@ from shapely.geometry import box from ldn.utils import ( - SENSOR, WGS84, LdnError, get_public_url_base, @@ -308,6 +307,7 @@ def build_pipeline_components( overwrite: bool, collection_url_root: str, s3_client: BaseClient, + sensor: str, ) -> tuple[PrefixedS3ItemPath, StacCreator, AwsDsCogWriter] | None: """Build shared pipeline components for GeoMAD and classify tasks. @@ -319,7 +319,7 @@ def build_pipeline_components( key_prefix=source_coop_prefix, prefix=owner, bucket=bucket, - sensor=SENSOR, + sensor=sensor, dataset_id=dataset_id, version=version, time=year, diff --git a/ldn/tests/test_geomad_cli.py b/ldn/tests/test_geomad_cli.py index 63264444..8cc6d1dd 100644 --- a/ldn/tests/test_geomad_cli.py +++ b/ldn/tests/test_geomad_cli.py @@ -132,7 +132,6 @@ def test_geomad_run_and_skip(bucket_env, mock_s3, runner, stac_key): "pacific", "--integration-test", "--overwrite", - "--no-single-region", ], ) assert result.exit_code == 0, result.output @@ -155,7 +154,6 @@ def test_geomad_run_and_skip(bucket_env, mock_s3, runner, stac_key): "--region", "pacific", "--integration-test", - "--no-single-region", ], ) assert result.exit_code == 0, result.output diff --git a/ldn/tests/test_raster.py b/ldn/tests/test_raster.py index b40de7d6..5e1e7225 100644 --- a/ldn/tests/test_raster.py +++ b/ldn/tests/test_raster.py @@ -552,6 +552,7 @@ def test_build_pipeline_components_uses_public_url_base_for_full_path_prefix(moc overwrite=True, collection_url_root="https://example.com/collections", s3_client=MagicMock(), + sensor="ls", ) mock_get_url.assert_called_once_with("dep-public-staging") itempath, *_ = result @@ -579,6 +580,7 @@ def test_build_pipeline_components_itempath_key_prefix(mock_aws, bucket, source_ overwrite=True, collection_url_root="https://example.com/collections", s3_client=MagicMock(), + sensor="ls", ) assert result is not None itempath, *_ = result @@ -599,6 +601,7 @@ def test_build_pipeline_components_returns_none_when_exists(mock_aws): overwrite=False, collection_url_root="https://example.com/collections", s3_client=MagicMock(), + sensor="ls", ) assert result is None @@ -617,5 +620,6 @@ def test_build_pipeline_components_proceeds_when_exists_and_overwrite(mock_aws): overwrite=True, collection_url_root="https://example.com/collections", s3_client=MagicMock(), + sensor="ls", ) assert result is not None diff --git a/ldn/training_data.py b/ldn/training_data.py index 39e97391..0667ed70 100644 --- a/ldn/training_data.py +++ b/ldn/training_data.py @@ -42,6 +42,7 @@ CLASS_ATTR, GEOMAD_DATASET_ID, GEOMAD_VERSION, + LULC_DATASET_ID, SENSOR, TRAINING_DATA_VERSION, TRAINING_DATA_YEAR, @@ -50,6 +51,8 @@ dataset_prefix, get_analysis_epsg, get_env_var, + get_public_url_base, + get_stac_geoparquet_key, get_stac_geoparquet_url, is_bucket_source_coop, owner_for_region, @@ -570,13 +573,13 @@ def filter_outliers(samples: gpd.GeoDataFrame, cap: float = 0.05) -> gpd.GeoData return samples -def _upload_dataframe_csv_to_s3(df, bucket: str, path: str) -> str: +def _upload_dataframe_csv_to_s3(df, bucket: str, key: str): """Upload a dataframe as CSV to S3 and return the S3 URI. Args: df: DataFrame to upload. bucket: S3 bucket name. - path: Key path within the bucket. + key: Key path within the bucket. Returns: S3 URI of the uploaded file. @@ -584,14 +587,13 @@ def _upload_dataframe_csv_to_s3(df, bucket: str, path: str) -> str: csv_buffer = io.StringIO() df.to_csv(csv_buffer, index=False) - key = f"{path}" s3_client.put_object( Bucket=bucket, Key=key, Body=csv_buffer.getvalue(), ContentType="text/csv", ) - return f"s3://{bucket}/{key}" # TODO: Use utils functions for S3 URI formatting. + logger.info(f"Uploaded CSV to {get_public_url_base(bucket)}/{key}") # Dep tools utils have mask_to_gadm() which would be helpful, but I want to buffer gadm before masking. @@ -638,13 +640,11 @@ def get_buffered_country( def get_tile_year_geomad_dem_indices( tile_id: str, year: str, - region: Literal["pacific", "non-pacific"], country_wgs84_buffered: gpd.GeoDataFrame, analysis_crs: Literal["EPSG:3832", "EPSG:6933"], product_owner: str, bucket: str, geomad_version: str, - single_region: bool, sensor: str, ) -> xr.Dataset: """Load GeoMAD + DEM features for a tile, clipped to buffered country. @@ -669,13 +669,11 @@ def get_tile_year_geomad_dem_indices( merged = search_and_load_geomad_indices_dem( tile_id=tile_id, year=year, - region=region, analysis_crs=analysis_crs, geopolygon=country_wgs84_buffered, product_owner=product_owner, geomad_version=geomad_version, bucket=bucket, - single_region=single_region, sensor=sensor, ) @@ -706,7 +704,6 @@ def make_training_data( file_prefix: str, n: int, min_sample_per_class_n: int, - single_region: bool, sensor: str, ): """Generate training data for a single tile and upload to S3 as CSV. @@ -765,17 +762,14 @@ def make_training_data( logger.info("Loading GeoMAD, DEM, and indices") owner = owner_for_region(region, product_owner) - # TODO: Use build_prefix() here. geomad_dem_indices = get_tile_year_geomad_dem_indices( tile_id, year, - region=region, country_wgs84_buffered=country_wgs84_buffered, analysis_crs=analysis_crs, product_owner=owner, bucket=geomad_bucket, geomad_version=geomad_version, - single_region=single_region, sensor=sensor, ) geobox = geomad_dem_indices.odc.geobox @@ -807,8 +801,8 @@ def make_training_data( samples.to_csv(f"{out_fname_local}.csv", index=False) logger.info(f"Saved training data to {out_fname_local}") - s3_uri = _upload_dataframe_csv_to_s3(samples, output_bucket, f"{file_prefix}.csv") - logger.info(f"Uploaded training data to {s3_uri}") + out_fname_s3 = f"{dataset_prefix(owner, sensor, LULC_DATASET_ID)}/{file_prefix}.csv" + _upload_dataframe_csv_to_s3(samples, output_bucket, f"{out_fname_s3}") @cli_training_app.command() @@ -835,6 +829,7 @@ def generate_training_data( n: int = typer.Option(2100, help="Total number of sample points"), min_sample_per_class_n: int = typer.Option(300, help="Minimum samples per class"), overwrite: bool = typer.Option(False, help="Whether to overwrite existing data in S3"), + # TODO: product_owner and single_region need to be seperate for geomad and output buckets. product_owner: str | None = typer.Option(None, help="Override the default product owner"), single_region: bool = typer.Option( ..., @@ -904,7 +899,6 @@ def generate_training_data( file_prefix, n, min_sample_per_class_n, - single_region, sensor, ) @@ -933,13 +927,11 @@ def make_geomad_item_id( def search_and_load_geomad_indices_dem( tile_id: str, year: str, - region: Literal["pacific", "non-pacific"], analysis_crs: Literal["EPSG:3832", "EPSG:6933"], geopolygon: gpd.GeoDataFrame, product_owner: str, bucket: str, geomad_version: str, - single_region: bool, sensor: str, ) -> xr.Dataset: """Search, load, scale, and merge GeoMAD bands, spectral indices, and DEM terrain for a tile. @@ -959,12 +951,11 @@ def search_and_load_geomad_indices_dem( Merged dataset with GeoMAD bands, spectral indices, elevation, slope, and aspect, clipped to the tile proj:bbox. """ - # geomad_stac_geoparquet_key = get_stac_geoparquet_key( - # bucket, single_region, product_owner, sensor, "geomad", geomad_version - # ) - geomad_stac_geoparquet_key = "" # TODO: Fix this. + geomad_stac_geoparquet_key = get_stac_geoparquet_key( + bucket, product_owner, sensor, GEOMAD_DATASET_ID, geomad_version + ) geomad_stac_geoparquet_url = get_stac_geoparquet_url(bucket, geomad_stac_geoparquet_key) - item_id = make_geomad_item_id(tile_id, sensor, year, product_owner=product_owner) + item_id = make_geomad_item_id(tile_id, sensor, year, product_owner) logger.info(f"Searching for GeoMAD item for tile {tile_id} and year {year}.") if GEOMAD_VERSION != geomad_version: @@ -975,6 +966,7 @@ def search_and_load_geomad_indices_dem( geomad_items = search_sync( geomad_stac_geoparquet_url, ids=item_id, + # max_items=1 ) geomad_items = [Item.from_dict(doc) for doc in geomad_items] geomad_items_n = len(geomad_items) diff --git a/notebooks/1_Train_Model.ipynb b/notebooks/1_Train_Model.ipynb index 97407d28..6ffbfcd3 100644 --- a/notebooks/1_Train_Model.ipynb +++ b/notebooks/1_Train_Model.ipynb @@ -74,7 +74,7 @@ "from ldn.utils import (\n", " CLASS_ATTR,\n", " MODEL_VERSION,\n", - " # TRAINING_DATA_VERSION,\n", + " TRAINING_DATA_VERSION,\n", " TRAINING_DATA_YEAR,\n", " WGS84,\n", " get_env_var,\n", @@ -84,7 +84,7 @@ "year = TRAINING_DATA_YEAR\n", "\n", "BUCKET = get_env_var(\"BUCKET\")\n", - "TRAINING_DATA_VERSION = \"0-0-4\" # Override. TODO: Run 0-0-9 for all pacific training tiles.\n", + "print(f\"BUCKET: {BUCKET}\")\n", "\n", "TRAINING_TILES = PACIFIC_TRAINING_TILES\n", "\n", @@ -104,17 +104,30 @@ "# Gather all of the CSVs of training points for many AOIs.\n", "# paths = f\"training_data/{TRAINING_DATA_VERSION}/{tile_id}/{year}/samples_*.csv\"\n", "# files = glob.glob(paths)\n", - "from ldn.aws import s3_client\n", - "from ldn.utils import parse_tile_id\n", + "import os\n", + "\n", + "import boto3\n", + "\n", + "# from ldn.aws import s3_client\n", + "from ldn.utils import LULC_DATASET_ID, dataset_prefix, owner_for_region, parse_tile_id\n", + "\n", + "aws_profile = os.environ.get(\"AWS_PROFILE\")\n", + "# %env AWS_PROFILE\n", + "print(f\"Using AWS profile: {aws_profile}\")\n", + "\n", + "aws_session = boto3.Session(profile_name=aws_profile)\n", + "print(f\"Using AWS profile: {aws_session.profile_name}, region: {aws_session.region_name}\")\n", + "s3_client = aws_session.client(\"s3\", region_name=aws_session.region_name)\n", "\n", "# List all training CSVs for this version\n", - "resp = s3_client.list_objects_v2(Bucket=BUCKET, Prefix=f\"training_data/{TRAINING_DATA_VERSION}/\")\n", + "owner = owner_for_region(region)\n", + "_dataset_prefix = dataset_prefix(owner, \"ls\", LULC_DATASET_ID)\n", + "resp = s3_client.list_objects_v2(Bucket=BUCKET, Prefix=f\"{_dataset_prefix}/training_data/{TRAINING_DATA_VERSION}/\")\n", "# This gets all regions (Pacific and non-Pacific). Can easily split by region here.\n", "keys = [obj[\"Key\"] for obj in resp.get(\"Contents\", []) if obj[\"Key\"].endswith(\"samples.csv\")]\n", "num_files_found = len(keys)\n", "print(f\"Found {num_files_found} CSV files\")\n", "\n", - "\n", "# Filter to just the training tiles.\n", "filtered_keys = []\n", "for id, region, country_dict in TRAINING_TILES:\n", @@ -127,7 +140,6 @@ "print(f\"Filtered to {num_files_found} CSV files for training tiles\")\n", "print(filtered_keys)\n", "\n", - "\n", "expected_num_files = len(TRAINING_TILES)\n", "assert num_files_found == expected_num_files, (\n", " f\"There should be {expected_num_files} training CSV files found in S3, not {num_files_found}.\"\n", @@ -162,14 +174,6 @@ "# Tiles are often missing one or more classes e.g. Singapore is missing 4:Wetland." ] }, - { - "cell_type": "code", - "execution_count": null, - "id": "e95dac89", - "metadata": {}, - "outputs": [], - "source": [] - }, { "cell_type": "code", "execution_count": null, @@ -524,18 +528,19 @@ "source": [ "from pathlib import Path\n", "\n", - "from ldn.utils import MODEL_VERSION\n", + "from ldn.utils import MODEL_VERSION, get_pathstyle_url_base\n", "\n", "s3_key = f\"models/{MODEL_VERSION}/{region}/{year}/lulc_random_forest_model_{region}_{year}.joblib\"\n", "\n", - "joblib_path = Path(f\"../../ldn/{s3_key}\")\n", + "joblib_path = Path(f\"../ldn/{s3_key}\")\n", "joblib_path.parent.mkdir(parents=True, exist_ok=True)\n", "joblib.dump(model, joblib_path)\n", "\n", "# Write the model to S3 so it can be loaded in prediction step.\n", + "s3_key = f\"{_dataset_prefix}/{s3_key}\"\n", "s3_client.upload_file(str(joblib_path), BUCKET, s3_key)\n", "\n", - "s3_url = f\"https://s3.us-west-2.amazonaws.com/{BUCKET}/{s3_key}\"\n", + "s3_url = f\"{get_pathstyle_url_base(BUCKET)}/{s3_key}\"\n", "print(f\"Uploaded model to {s3_url}\")\n", "\n", "# # Load the model\n", diff --git a/templates/ls-geomad-annual.yaml b/templates/ls-geomad-annual.yaml index b6470c69..2e8e7ec2 100644 --- a/templates/ls-geomad-annual.yaml +++ b/templates/ls-geomad-annual.yaml @@ -43,9 +43,6 @@ spec: - name: region # Just run Pacific in DEP Argo instance. value: "pacific" # "all", "pacific", or "non-pacific" - - name: single-region - value: "--single-region" # For single region products e.g. DEP (staging and prod). - # value: "false" # For places with both regions e.g. Source Coop. - name: years # Need to run in 8 years blocks because print-tasks output is too large. # value: "2000-2008" @@ -97,8 +94,8 @@ spec: # provision a small node first (c6.xlarge/large), then process geomad can't fit # and has to wait for another node to come up. Not critical, just saves a few minutes # of waiting between tasks. - memory: 16Gi - cpu: 10 + memory: 10Gi + cpu: 3 command: [sh, -c] args: - | @@ -159,7 +156,6 @@ spec: - "{{ inputs.parameters.year }}" - --region - "{{ inputs.parameters.region }}" - - "{{ workflow.parameters.single-region }}" - --version - "{{ workflow.parameters.version }}" - --collection-url-root diff --git a/templates/ls-lulc-classify-annual.yaml b/templates/ls-lulc-classify-annual.yaml index 9080580c..9fc5525b 100644 --- a/templates/ls-lulc-classify-annual.yaml +++ b/templates/ls-lulc-classify-annual.yaml @@ -5,8 +5,8 @@ metadata: namespace: argo spec: entrypoint: workflow-entrypoint - # serviceAccountName: public-bucket-writer # DEP Staging - serviceAccountName: open-data-bucket-writer # DEP Prod + serviceAccountName: public-bucket-writer # DEP Staging + # serviceAccountName: open-data-bucket-writer # DEP Prod ttlStrategy: secondsAfterCompletion: 300 podGC: @@ -21,13 +21,13 @@ spec: # TODO: Need to disactivate all dep-specific tags to run non-Pacific processes in another Argo instance. nodeSelector: karpenter.sh/capacity-type: "spot" # Just to be explicit. In dep-terraform, the data-processing nodepool can only pick spot instances - # dep/workload: data-processing # DEP Staging - dep/workload: data-processing-public # DEP Prod + dep/workload: data-processing # DEP Staging + # dep/workload: data-processing-public # DEP Prod tolerations: - key: dep/workload operator: Equal - # value: data-processing # DEP Staging - value: data-processing-public # DEP Prod + value: data-processing # DEP Staging + # value: data-processing-public # DEP Prod effect: NoSchedule arguments: parameters: @@ -36,27 +36,29 @@ spec: - name: image-tag value: "latest" - name: lulc-version - value: "0-0-4" + value: "0-0-9" - name: geomad-version - value: "0-3-0" + value: "0-2-1" - name: bucket - # value: "dep-public-staging" # DEP Staging - value: "dep-public-data" # DEP Prod + value: "dep-public-staging" # DEP Staging + # value: "dep-public-data" # DEP Prod - name: region # Just run Pacific in DEP Argo instance. value: "pacific" # "all", "pacific", or "non-pacific" - name: years - # TODO: Remove workaround of year subset blocks once the JSON issue is resolved. - # value: "2000-2008" - # value: "2009-2016" - # value: "2017-2025" - # value: "2000-2025" - value: "2025" # Training data is 2020 but predicting on 2025 because there are less gaps than 2020. + value: "2000,2025" + # Training data is 2020 but predicting on 2025 because there are less gaps than 2020. + - name: model-path + value: "https://dep-public-staging.s3.us-west-2.amazonaws.com/dep_ls_lulc/models/0-0-9/pacific/2020/lulc_random_forest_model_pacific_2020.joblib" + - name: collection-url-root + value: "https://stac.staging.digitalearthpacific.io/collections" # DEP Staging + # value: "https://stac.digitalearthpacific.org/collections" # DEP Prod + # value: None # For Source Coop, let collection-url-root default - name: overwrite - # value: "--overwrite" # Can be "--overwrite" or "--no-overwrite" - value: "--no-overwrite" # Can be "--overwrite" or "--no-overwrite" + # value: "--overwrite" + value: "--no-overwrite" - name: decimated - value: "--no-decimated" # Can be "--decimated" or "--no-decimated" + value: "--no-decimated" templates: - name: workflow-entrypoint dag: @@ -92,14 +94,19 @@ spec: # of waiting between tasks. requests: memory: 8Gi - cpu: 4 - limits: - memory: 12Gi - cpu: 6 + cpu: 3 command: [sh, -c] args: - | - ldn print-tasks --years "{{ workflow.parameters.years }}" --region "{{ workflow.parameters.region }}" --lulc-version "{{ workflow.parameters.lulc-version }}" --dataset lulc {{ workflow.parameters.overwrite }} > /tmp/tasks.json + ldn print-tasks \ + --years "{{ workflow.parameters.years }}" \ + --region "{{ workflow.parameters.region }}" \ + --lulc-version "{{ workflow.parameters.lulc-version }}" \ + --geomad-version "{{ workflow.parameters.geomad-version }}" \ + --dataset lulc \ + {{ workflow.parameters.overwrite }} \ + --bucket "{{ workflow.parameters.bucket }}" \ + > /tmp/tasks.json outputs: parameters: - name: tasks @@ -120,10 +127,10 @@ spec: resources: requests: memory: 8Gi - cpu: 4 + cpu: 3 limits: memory: 12Gi - cpu: 6 + cpu: 4 command: [ldn] args: - lulc @@ -138,16 +145,19 @@ spec: - "{{ workflow.parameters.lulc-version }}" - --geomad-version - "{{ workflow.parameters.geomad-version }}" - # Let model_path default - # - --model-path - # - "https://s3.us-west-2.amazonaws.com/data.ldn.auspatious.com/models/0-0-4/pacific/2020/lulc_random_forest_model_pacific_2020.joblib" - - "{{ workflow.parameters.decimated }}" - - "{{ workflow.parameters.overwrite }}" + - --model-path + - "{{ workflow.parameters.model-path }}" + - --bucket + - "{{ workflow.parameters.bucket }}" + - --collection-url-root + - "{{ workflow.parameters.collection-url-root }}" - --memory-limit - - "10GB" + - "6GB" - --n-workers - "2" - --threads-per-worker - "4" - --xy-chunk-size - "1024" + - "{{ workflow.parameters.decimated }}" + - "{{ workflow.parameters.overwrite }}"