Skip to content
Open
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
25 changes: 19 additions & 6 deletions argopy/data_fetchers/erddap_data.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,6 @@
from .proto import ArgoDataFetcherProto
from .erddap_data_processors import pre_process, quote_string_constraints


log = logging.getLogger("argopy.erddap.data")

access_points = ["wmo", "box"]
Expand All @@ -48,6 +47,7 @@ class ErddapArgoDataFetcher(ArgoDataFetcherProto):
This class is a prototype not meant to be instantiated directly

"""

data_source = "erddap"

###
Expand Down Expand Up @@ -146,8 +146,8 @@ def __init__( # noqa: C901
self.fs = kwargs["fs"] if "fs" in kwargs else httpstore(**self.store_opts)

self.parallelize, self.parallel_method = PARALLEL_SETUP(parallel)
if self.parallelize and self.parallel_method == 'thread':
self.parallel_method = 'erddap' # Use our custom filestore
if self.parallelize and self.parallel_method == "thread":
self.parallel_method = "erddap" # Use our custom filestore
self.progress = progress
self.chunks = chunks
self.chunks_maxsize = chunks_maxsize
Expand Down Expand Up @@ -238,6 +238,17 @@ def __init__( # noqa: C901
% (v, ", ".join(self._bgc_vlist_avail))
)

if (
self._bgc_vlist_params is not None
and self._bgc_vlist_measured is not None
):
for v in self._bgc_vlist_measured:
if v not in self._bgc_vlist_params:
raise ValueError(
"'%s' can't be in the 'measured' list if it is not in the 'params' list. Please check your arguments."
% v
)

def __repr__(self):
summary = ["<datafetcher.erddap>"]
summary.append(self._repr_data_source)
Expand Down Expand Up @@ -358,7 +369,7 @@ def _minimal_vlist(self):
]
[vlist.append(p) for p in plist]

# Core/Deep variables:
# BGC/Core/Deep variables:
plist = [p.lower() for p in list_core_parameters()]
[vlist.append(p) for p in plist]
[vlist.append(p + "_qc") for p in plist]
Expand Down Expand Up @@ -520,7 +531,9 @@ def getNfromncHeader(url):
if "Your query produced no matching results. (nRows = 0)" in ncHeader:
return 0
else:
lines = [line for line in ncHeader.splitlines() if "row = " in line][0]
lines = [
line for line in ncHeader.splitlines() if "row = " in line
][0]
return int(lines.split("=")[1].split(";")[0])
except Exception:
raise ErddapServerError(
Expand Down Expand Up @@ -943,6 +956,6 @@ def uri(self):
except DataNotFound:
log.debug("This box fetcher will contain no data")
except ValueError as e:
if 'not available for this access point' in str(e):
if "not available for this access point" in str(e):
log.debug("This box fetcher does not contained required data")
return urls
212 changes: 210 additions & 2 deletions argopy/data_fetchers/gdac_data.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,13 +9,16 @@
import pandas as pd
import xarray as xr
from abc import abstractmethod
from typing import Union
import warnings
import getpass
import logging
from typing import Literal

from ..utils.format import argo_split_path
from ..options import OPTIONS, check_gdac_option, PARALLEL_SETUP
from ..utils.lists import list_bgc_s_variables, list_core_parameters
from ..utils import is_list_of_strings, to_list
from ..errors import DataNotFound
from ..stores import ArgoIndex, has_distributed, distributed
from .proto import ArgoDataFetcherProto
Expand Down Expand Up @@ -64,6 +67,8 @@ def __init__(
dimension: Literal["point", "profile"] = "point",
errors: str = "raise",
api_timeout: int = 0,
params: Union[str, list] = "all",
measured: Union[str, list] = None,
**kwargs
):
"""Init fetcher
Expand Down Expand Up @@ -94,6 +99,14 @@ def __init__(
Show a progress bar or not when fetching data.
api_timeout: int (optional)
Server request time out in seconds. Set to OPTIONS['api_timeout'] by default.
params: Union[str, list] (optional, default='all')
List of BGC essential variables to retrieve, i.e. that will be in the output :class:`xr.DataSet``.
By default, this is set to ``all``, i.e. any variable found in at least of the profile in the data
selection will be included in the output.
measured: Union[str, list] (optional, default=None)
List of BGC essential variables that can't be NaN. If set to 'all', this is an easy way to reduce the size of the
:class:`xr.DataSet`` to points where all variables have been measured. Otherwise, provide a simple list of
variables.

Other parameters
----------------
Expand All @@ -108,6 +121,8 @@ def __init__(
self.errors = errors
self.dimension = dimension

self.init(**kwargs)

# Validate server, raise GdacPathError if not valid.
check_gdac_option(self.server, errors="raise")

Expand All @@ -123,8 +138,81 @@ def __init__(
cachedir=cachedir,
timeout=self.timeout,
)

self.fs = self.indexfs.fs["src"] # Reuse the appropriate file system

if self.dataset_id in ["bgc", "bgc-s"]:
self._bgc_vlist_gdac = [v.lower() for v in list_bgc_s_variables()]
# Handle the 'params' argument:
self._bgc_params = to_list(params)
if isinstance(params, str):
if params == "all":
params = self._bgc_vlist_avail
else:
params = to_list(params)
elif params is None:
raise ValueError()
elif params[0] == "all":
params = self._bgc_vlist_avail
elif not is_list_of_strings(params):
raise ValueError("'params' argument must be a list of strings")
# raise ValueError("'params' argument must be a list of strings (possibly with a * wildcard)")
self._bgc_vlist_params = [p.upper() for p in params]
# self._bgc_vlist_params = self._bgc_handle_wildcard(self._bgc_vlist_params)

for v in self._bgc_vlist_params:
if v not in self._bgc_vlist_avail:
raise ValueError(
"'%s' not available for this access point. The 'params' argument must have values in [%s]"
% (v, ",".join(self._bgc_vlist_avail))
)

for p in list_core_parameters():
if p not in self._bgc_vlist_params:
self._bgc_vlist_params.append(p)

if (
self.user_mode in ["standard", "research"]
and "CDOM" in self._bgc_vlist_params
):
self._bgc_vlist_params.remove("CDOM")
log.warning(
"CDOM was requested but was removed from the fetcher because executed in '%s' user mode"
% self.user_mode
)

# Handle the 'measured' criteria for BGC requests:
self._bgc_measured = to_list(measured)
if isinstance(measured, str):
if measured == "all":
measured = self._bgc_vlist_params
else:
measured = to_list(measured)
elif self._bgc_measured[0] is None:
measured = []
elif self._bgc_measured[0] == "all":
measured = self._bgc_vlist_params
elif not is_list_of_strings(self._bgc_measured):
raise ValueError("'measured' argument must be a list of strings")
# raise ValueError("'measured' argument must be a list of strings (possibly with a * wildcard)")
self._bgc_vlist_measured = [m.upper() for m in measured]
# self._bgc_vlist_measured = self._bgc_handle_wildcard(self._bgc_vlist_measured)

for v in self._bgc_vlist_measured:
if v not in self._bgc_vlist_avail:
raise ValueError(
"'%s' not available for this access point. The 'measured' argument must have values in [%s]"
% (v, ", ".join(self._bgc_vlist_avail))
)

if self._bgc_vlist_params is not None and self._bgc_vlist_measured is not None:
for v in self._bgc_vlist_measured:
if v not in self._bgc_vlist_params:
raise ValueError(
"'%s' can't be in the 'measured' list if it is not in the 'params' list. Please check your arguments."
% v
)

nrows = None
if "N_RECORDS" in kwargs:
nrows = kwargs["N_RECORDS"]
Expand All @@ -136,13 +224,16 @@ def __init__(
self.parallelize, self.parallel_method = PARALLEL_SETUP(parallel)
self.progress = progress

self.init(**kwargs)

def __repr__(self):
summary = ["<datafetcher.gdac>"]
summary.append(self._repr_data_source)
summary.append(self._repr_access_point)
summary.append(self._repr_server)
if self.dataset_id in ["bgc", "bgc-s"]:
summary.append("📗 Parameters: %s" % self._bgc_vlist_params)
summary.append(
"📕 BGC 'must be measured' parameters: %s" % self._bgc_vlist_measured
)
if hasattr(self.indexfs, "index"):
summary.append(
"📗 Index: %s (%i records)" % (self.indexfs.index_file, self.N_RECORDS)
Expand Down Expand Up @@ -227,6 +318,117 @@ def mono2multi(mono_path):
new_uri = list(set(new_uri))
return new_uri

@property
def _bgc_vlist_avail(self):
"""Return the list of the gdac BGC dataset available for this access point

Apply search criteria in the index, then retrieve the list of parameters
"""
if hasattr(self, "WMO"):
if hasattr(self, "CYC") and self.CYC is not None:
self.indexfs.query.wmo_cyc(self.WMO, self.CYC)
else:
self.indexfs.query.wmo(self.WMO)
elif hasattr(self, "BOX"):
if len(self.indexBOX) == 4:
self.indexfs.query.lon_lat(self.indexBOX)
else:
self.indexfs.query.box(self.indexBOX)

params = self.indexfs.read_params()

results = []
for p in params:
if p.lower() in self._bgc_vlist_gdac:
results.append(p)
# else:
# log.error(
# "Removed '%s' because it is not available on the erddap server (%s), but it should !"
# % (p, self._server)
# )

return results

@property
def _minimal_vlist(self):
"""Return the list of variables to retrieve measurements for"""
vlist = list()
if self.dataset_id == "phy":
plist = [
"data_mode",
"latitude",
"longitude",
"position_qc",
"time",
"time_qc",
"direction",
"platform_number",
"cycle_number",
"config_mission_number",
"vertical_sampling_scheme",
]
[vlist.append(p) for p in plist]

# BGC/Core/Deep variables:
plist = [p.lower() for p in list_core_parameters()]
[vlist.append(p) for p in plist]
[vlist.append(p + "_qc") for p in plist]
[vlist.append(p + "_adjusted") for p in plist]
[vlist.append(p + "_adjusted_qc") for p in plist]
[vlist.append(p + "_adjusted_error") for p in plist]

if self.dataset_id in ["bgc", "bgc-s"]:
plist = [
# "parameter_data_mode", # never !!!
"latitude",
"longitude",
"position_qc",
"time",
"time_qc",
"direction",
"platform_number",
"cycle_number",
"config_mission_number",
]
[vlist.append(p) for p in plist]

# Search in the profile index the list of parameters to load:
params = self._bgc_vlist_params # rq: include 'core' variables
# log.debug("erddap-bgc parameters to load: %s" % params)

for p in params:
vname = p.lower()
if self.user_mode in ["expert"]:
vlist.append("%s" % vname)
vlist.append("%s_qc" % vname)
vlist.append("%s_adjusted" % vname)
vlist.append("%s_adjusted_qc" % vname)
vlist.append("%s_adjusted_error" % vname)

elif self.user_mode in ["standard"]:
vlist.append("%s" % vname)
vlist.append("%s_qc" % vname)
vlist.append("%s_adjusted" % vname)
vlist.append("%s_adjusted_qc" % vname)
vlist.append("%s_adjusted_error" % vname)

elif self.user_mode in ["research"]:
vlist.append("%s_adjusted" % vname)
vlist.append("%s_adjusted_qc" % vname)
vlist.append("%s_adjusted_error" % vname)

# vlist.append("profile_%s_qc" % vname) # not in the database

if self.dataset_id == "ref":
plist = ["latitude", "longitude", "time", "platform_number", "cycle_number"]
[vlist.append(p) for p in plist]
plist = ["pres", "temp", "psal", "ptmp"]
[vlist.append(p) for p in plist]

vlist.sort()
vlist = [p.upper() for p in vlist]
return vlist

@property
def cachepath(self):
"""Return path to cache file(s) for this request
Expand Down Expand Up @@ -298,6 +500,8 @@ def to_xarray(
"access_point_opts": access_point_opts,
"pre_filter_points": self._post_filter_points,
"dimension": dimension,
"params_list": self._minimal_vlist,
"measured_params": self._bgc_vlist_measured if self.dataset_id in ["bgc", "bgc-s"] else None,
}

# Download and pre-process data:
Expand All @@ -309,6 +513,10 @@ def to_xarray(
"preprocess": pre_process_multiprof,
"preprocess_opts": preprocess_opts,
}
# ATTEMPT TO HANDLE BGC PARAMS SELECTION FOR GDAC
if self.dataset_id in ["bgc", "bgc-s"]:
opts["data_vars"] = self._bgc_vlist_params

if self.parallel_method in ["thread"]:
opts["method"] = "thread"
opts["open_dataset_opts"] = {"xr_opts": {"engine": "argo"}}
Expand Down
Loading
Loading