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
2 changes: 1 addition & 1 deletion pygeoapi/api/collection.py
Original file line number Diff line number Diff line change
Expand Up @@ -437,7 +437,7 @@ def gen_collection(api, request, dataset: str,
'type': FORMAT_TYPES[F_HTML],
'rel': 'data',
'title': title2,
'href': f'{api.get_collections_url()}/{dataset}/{qt}?f={F_HTML}' # noqa
'href': f'{api.get_collections_url()}/{dataset}/{qt}?f={F_HTML}' # noqa
}])

for key, value in get_dataset_formatters(config).items():
Expand Down
207 changes: 152 additions & 55 deletions pygeoapi/provider/mongo.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,6 @@
#
# =================================================================

from datetime import datetime
import logging

from pymongo import MongoClient
Expand All @@ -43,8 +42,7 @@


class MongoProvider(BaseProvider):
"""Generic provider for Mongodb.
"""
"""Generic provider for Mongodb."""

def __init__(self, provider_def):
"""
Expand All @@ -57,16 +55,18 @@ def __init__(self, provider_def):
"""
# this is dummy value never used in case of Mongo.
# Mongo id field is _id
provider_def.setdefault('id_field', '_id')
provider_def.setdefault("id_field", "_id")

super().__init__(provider_def)

LOGGER.info(f'Mongo source config: {self.data}')
LOGGER.info(f"Mongo source config: {self.data}")

dbclient = MongoClient(self.data)
self.featuredb = dbclient.get_default_database()
self.collection = provider_def['collection']
self.featuredb[self.collection].create_index([("geometry", GEOSPHERE)])
self.collection = provider_def["collection"]
self.featuredb[self.collection].create_index(
[("feature.geometry", GEOSPHERE)]
)
self.get_fields()

def get_fields(self):
Expand All @@ -81,7 +81,7 @@ def get_fields(self):
{"$project": {"properties": 1}},
{"$unwind": "$properties"},
{"$group": {"_id": "$properties", "count": {"$sum": 1}}},
{"$project": {"_id": 1}}
{"$project": {"_id": 1}},
]

result = list(self.featuredb[self.collection].aggregate(pipeline))
Expand All @@ -90,80 +90,176 @@ def get_fields(self):
# set the field type to 'string'.
# by operating without a schema, mongo can query any data type.
for i in result:
for key in result[0]['_id'].keys():
self._fields[key] = {'type': 'string'}
for key in result[0]["_id"].keys():
self._fields[key] = {"type": "string"}

return self._fields

def _get_feature_list(self, filterObj, sortList=[], skip=0, maxitems=1,
skip_geometry=False):
def _get_feature_list(
self, filterObj, sortList=[], skip=0, maxitems=1, skip_geometry=False
):
featurecursor = self.featuredb[self.collection].find(filterObj)

if sortList:
featurecursor = featurecursor.sort(sortList)

featurecursor.skip(skip)
featurecursor.limit(maxitems)
if maxitems > -1:
featurecursor.limit(maxitems)
featurelist = list(featurecursor)

features = []

for item in featurelist:
item['id'] = str(item.pop('_id'))
feature_id = str(item.pop("_id"))
geometry = item["feature"]["geometry"]
props = item["feature"]["properties"]

if skip_geometry:
item['geometry'] = None
geometry = None

feature = {
"type": "Feature",
"id": feature_id,
"geometry": geometry,
"properties": props,
}

features.append(feature)

return featurelist
return features

@crs_transform
def query(self, offset=0, limit=10, resulttype='results',
bbox=[], datetime_=None, properties=[], sortby=[],
select_properties=[], skip_geometry=False, q=None, **kwargs):
def query(
self,
offset=0,
limit=10,
resulttype="results",
bbox=[],
datetime_=None,
properties=[],
sortby=[],
select_properties=[],
skip_geometry=False,
q=None,
filterq=None,
**kwargs,
):
"""
query the provider

:returns: dict of 0..n GeoJSON features
"""

def cql2_to_mongo(node):
if node is None:
return

# GeoJson operator
if node.__class__.__name__ == "GeometryWithin":
field = node.lhs.name
geom = node.rhs.geometry

query_body = {field: {"$geoWithin": {"$geometry": geom}}}
return query_body

if node.__class__.__name__ == "GeometryIntersects":
field = node.lhs.name
geom = node.rhs.geometry

return {field: {"$geoIntersects": {"$geometry": geom}}}

if node.__class__.__name__ == "DistanceWithin":
field = node.lhs.name
geom = node.rhs.geometry
distance = node.distance
# units = node.units # mongo's default units are meters
# but with CQL we can pass different units
# and here we can recalculate them
return {
field: {
"$near": {
"$geometry": geom,
"$maxDistance": distance,
"$minDistance": 0,
}
}
}

# Logical operators
if node.__class__.__name__ == "And":
return {"$and":
[cql2_to_mongo(node.lhs),
cql2_to_mongo(node.rhs)]
}

if node.__class__.__name__ == "Or":
return {"$or":
[cql2_to_mongo(node.lhs),
cql2_to_mongo(node.rhs)]
}

# Comparison operators
if node.__class__.__name__ == "Equal":
field = node.lhs.name
value = node.rhs
return {f"{field}": value}

if node.__class__.__name__ == "GreaterEqual":
field = node.lhs.name
value = node.rhs
return {f"{field}": {"$gte": value}}

if node.__class__.__name__ == "LessEqual":
field = node.lhs.name
value = node.rhs
return {f"{field}": {"$lte": value}}

return

and_filter = []
cql_filters_parsed = cql2_to_mongo(filterq)
if cql_filters_parsed is not None:
and_filter.append(cql_filters_parsed)
limit = -1 # if there is CQL query return all elements

if len(bbox) == 4:
x, y, w, h = map(float, bbox)
and_filter.append(
{'geometry': {'$geoWithin': {'$box': [[x, y], [w, h]]}}})

# This parameter is not working yet!
# gte is not sufficient to check date range
if datetime_ is not None:
assert isinstance(datetime_, datetime)
and_filter.append({'properties.datetime': {'$gte': datetime_}})
and_filter.append({
"geometry": {
"$geoWithin": {
"$box": [[x, y], [w, h]]
}
}
})

for prop in properties:
and_filter.append({"properties."+prop[0]: {'$eq': prop[1]}})
and_filter.append({"properties." + prop[0]: {"$eq": prop[1]}})

filterobj = {'$and': and_filter} if and_filter else {}
filterobj = {"$and": and_filter} if and_filter else {}

sort_list = [("properties." + sort['property'],
ASCENDING if (sort['order'] == '+') else DESCENDING)
for sort in sortby]
sort_list = [
(
"properties." + sort["property"],
ASCENDING if (sort["order"] == "+") else DESCENDING,
)
for sort in sortby
]

feature_collection = {
'type': 'FeatureCollection',
'features': []
}
feature_collection = {"type": "FeatureCollection", "features": []}

if self.count or resulttype == 'hits':
matched = self.featuredb[self.collection].count_documents(
filterobj)
LOGGER.debug(f'Found {matched} result(s)')
feature_collection['numberMatched'] = matched

if resulttype == 'hits':
if resulttype == "hits":
return feature_collection

featurelist = self._get_feature_list(
filterobj, sortList=sort_list, skip=offset, maxitems=limit,
skip_geometry=skip_geometry
filterobj,
sortList=sort_list,
skip=offset,
maxitems=limit,
skip_geometry=skip_geometry,
)

feature_collection['features'] = featurelist
feature_collection['numberReturned'] = len(featurelist)
feature_collection["features"] = featurelist
feature_collection["numberReturned"] = len(featurelist)

return feature_collection

Expand All @@ -175,17 +271,16 @@ def get(self, identifier, **kwargs):
:param identifier: feature id
:returns: dict of single GeoJSON feature
"""
featurelist = self._get_feature_list({'_id': ObjectId(identifier)})
featurelist = self._get_feature_list({"_id": ObjectId(identifier)})
if featurelist:
return featurelist[0]
else:
err = f'item {identifier} not found'
err = f"item {identifier} not found"
LOGGER.error(err)
raise ProviderItemNotFoundError(err)

def create(self, new_feature):
"""Create a new feature
"""
"""Create a new feature"""
self.featuredb[self.collection].insert_one(new_feature)

def update(self, identifier, updated_feature):
Expand All @@ -194,14 +289,16 @@ def update(self, identifier, updated_feature):
:param identifier: feature id
:param new_feature: new GeoJSON feature dictionary
"""
data = {k: v for k, v in updated_feature.items() if k != 'id'}
data = {k: v for k, v in updated_feature.items() if k != "id"}
self.featuredb[self.collection].update_one(
{'_id': ObjectId(identifier)}, {"$set": data})
{"_id": ObjectId(identifier)}, {"$set": data}
)

def delete(self, identifier):
"""Deletes an existing feature

:param identifier: feature id
"""
self.featuredb[self.collection].delete_one(
{'_id': ObjectId(identifier)})
{"_id": ObjectId(identifier)}
)