binary wire format with flatbuffers (#509)

* first flatbuffer schema

* do not lint auto-generated files

* add flatbuffers package

* add flatbuffer module

* wire up /data/X/T route

* use flatbuffers for matrix data fetc

* clarity and comments

* add flatbuffer layout route

* clean up obsolete code

* fix tests

* move flake8 config to setup.cfg

* add comments

* lint

* rework layout routes for fbs

* add more type support to fbs

* lint

* add flatbuffer support for annotations

* function name improvements

* fix botched merge with master

* remove unused import

* route cleanup for flatbuffers

* rename function for clarity

* add missing globals to Jest tests

* fix client JS tests

* fix routes for Python tests

* comments for clarity

* non-finite floating point hardening

* more non-finite number handling

* lint

* fix tests for summarizeAnnotations

* harden diffexp calculation against FP errors

* cleanup unused code

* lint

* add encoding tests for flatbuffers

* application type specified as strings

* fix spelling error

* improve variable names

* add note about documentation gap

* rename FBS DataFrame to Matrix
This commit is contained in:
Bruce Martin
2019-01-09 14:26:05 -08:00
committed by GitHub
parent 42e25a1a1f
commit b90447c387
41 changed files with 2317 additions and 266 deletions
+14
View File
@@ -66,6 +66,11 @@ class CXGDriver(metaclass=ABCMeta):
"""
pass
@abstractmethod
def annotation_to_fbs_matrix(self, axis, field=None):
""" Same as annotation(), except returns a flatbuffer, and does not support filtering. """
pass
@abstractmethod
def data_frame(self, filter, axis):
"""
@@ -79,6 +84,10 @@ class CXGDriver(metaclass=ABCMeta):
"""
pass
@abstractmethod
def data_frame_to_fbs_matrix(self, filter, axis):
pass
@abstractmethod
def diffexp_topN(self, obsFilter1, obsFilter2, top_n=None, interactive_limit=None):
"""
@@ -104,3 +113,8 @@ class CXGDriver(metaclass=ABCMeta):
:return: [cellid, x, y, ...]
"""
pass
@abstractmethod
def layout_to_fbs_matrix(self, filter):
""" same as layout, except returns a flatbuffer """
pass
+69 -39
View File
@@ -9,7 +9,6 @@ from werkzeug.datastructures import ImmutableMultiDict
from server.app.util.constants import (
Axis,
DiffExpMode,
JSON_MIMETYPE,
JSON_NaN_to_num_warning_msg,
)
from server.app.util.filter import parse_filter, QueryStringError
@@ -194,11 +193,21 @@ class AnnotationsObsAPI(Resource):
)
def get(self):
fields = request.args.getlist("annotation-name", None)
preferred_mimetype = request.accept_mimetypes.best_match(
["application/json", "application/octet-stream"],
"application/json"
)
try:
annotation_response = current_app.data.annotation({}, "obs", fields)
return make_response(
annotation_response, HTTPStatus.OK, {"Content-Type": JSON_MIMETYPE}
)
if preferred_mimetype == "application/json":
return make_response(
current_app.data.annotation({}, "obs", fields), HTTPStatus.OK, {"Content-Type": "application/json"}
)
elif preferred_mimetype == "application/octet-stream":
return make_response(current_app.data.annotation_to_fbs_matrix("obs", fields),
HTTPStatus.OK,
{"Content-Type": "application/octet-stream"})
else:
return make_response(f"Unsupported MIME type '{request.accept_mimetypes}'", HTTPStatus.NOT_ACCEPTABLE)
except KeyError:
return make_response(f"Error bad key in {fields}", HTTPStatus.BAD_REQUEST)
except JSONEncodingValueError as e:
@@ -254,7 +263,7 @@ class AnnotationsObsAPI(Resource):
request.get_json()["filter"], "obs", fields
)
return make_response(
annotation_response, HTTPStatus.OK, {"Content-Type": JSON_MIMETYPE}
annotation_response, HTTPStatus.OK, {"Content-Type": "application/json"}
)
except KeyError:
return make_response(f"Error bad key in {fields}", HTTPStatus.BAD_REQUEST)
@@ -304,11 +313,21 @@ class AnnotationsVarAPI(Resource):
)
def get(self):
fields = request.args.getlist("annotation-name", None)
preferred_mimetype = request.accept_mimetypes.best_match(
["application/json", "application/octet-stream"],
"application/json"
)
try:
annotation_response = current_app.data.annotation({}, "var", fields)
return make_response(
annotation_response, HTTPStatus.OK, {"Content-Type": JSON_MIMETYPE}
)
if preferred_mimetype == "application/json":
return make_response(current_app.data.annotation({}, "var", fields),
HTTPStatus.OK,
{"Content-Type": "application/json"})
elif preferred_mimetype == "application/octet-stream":
return make_response(current_app.data.annotation_to_fbs_matrix("var", fields),
HTTPStatus.OK,
{"Content-Type": "application/octet-stream"})
else:
return make_response(f"Unsupported MIME type '{request.accept_mimetypes}'", HTTPStatus.NOT_ACCEPTABLE)
except KeyError:
return make_response(f"Error bad key in {fields}", HTTPStatus.BAD_REQUEST)
except JSONEncodingValueError as e:
@@ -364,7 +383,7 @@ class AnnotationsVarAPI(Resource):
request.get_json()["filter"], "var", fields
)
return make_response(
annotation_response, HTTPStatus.OK, {"Content-Type": JSON_MIMETYPE}
annotation_response, HTTPStatus.OK, {"Content-Type": "application/json"}
)
except KeyError:
return make_response(f"Error bad key in {fields}", HTTPStatus.BAD_REQUEST)
@@ -437,7 +456,7 @@ class DataObsAPI(Resource):
return make_response(
current_app.data.data_frame(filter_, axis=Axis.OBS),
HTTPStatus.OK,
{"Content-Type": JSON_MIMETYPE},
{"Content-Type": "application/json"},
)
except FilterError as e:
return make_response(e.message, HTTPStatus.BAD_REQUEST)
@@ -495,7 +514,7 @@ class DataObsAPI(Resource):
)
),
HTTPStatus.OK,
{"Content-Type": JSON_MIMETYPE},
{"Content-Type": "application/json"},
)
except FilterError as e:
return make_response(e.message, HTTPStatus.BAD_REQUEST)
@@ -564,7 +583,7 @@ class DataVarAPI(Resource):
return make_response(
current_app.data.data_frame(filter_, axis=Axis.VAR),
HTTPStatus.OK,
{"Content-Type": JSON_MIMETYPE},
{"Content-Type": "application/json"},
)
except FilterError as e:
return make_response(e.message, HTTPStatus.BAD_REQUEST)
@@ -603,28 +622,30 @@ class DataVarAPI(Resource):
}
)
def put(self):
if not request.accept_mimetypes.best_match(["application/json", "text/csv"]):
return make_response(
f"Unsupported MIME type '{request.accept_mimetypes}'",
HTTPStatus.NOT_ACCEPTABLE,
)
# TODO support CSV
preferred_mimetype = request.accept_mimetypes.best_match(
["application/json", "application/octet-stream"],
"application/json"
)
try:
get_mime_type(
acceptable_types=["application/json"], header=request.accept_mimetypes
)
except MimeTypeError as e:
return make_response(e.message, HTTPStatus.NOT_ACCEPTABLE)
try:
return make_response(
(
current_app.data.data_frame(
if preferred_mimetype == "application/json":
return make_response(
(
current_app.data.data_frame(
request.get_json()["filter"], axis=Axis.VAR
)
),
HTTPStatus.OK,
{"Content-Type": "application/json"},
)
elif preferred_mimetype == "application/octet-stream":
return make_response(
current_app.data.data_frame_to_fbs_matrix(
request.get_json()["filter"], axis=Axis.VAR
)
),
HTTPStatus.OK,
{"Content-Type": JSON_MIMETYPE},
)
),
HTTPStatus.OK,
{"Content-Type": "application/octet-stream"})
else:
return make_response(f"Unsupported MIME type '{request.accept_mimetypes}'", HTTPStatus.NOT_ACCEPTABLE)
except FilterError as e:
return make_response(e.message, HTTPStatus.BAD_REQUEST)
except JSONEncodingValueError as e:
@@ -747,7 +768,7 @@ class DiffExpObsAPI(Resource):
current_app.data.features["diffexp"]["interactiveLimit"],
)
return make_response(
diffexp, HTTPStatus.OK, {"Content-Type": JSON_MIMETYPE}
diffexp, HTTPStatus.OK, {"Content-Type": "application/json"}
)
except (ValueError, FilterError) as e:
return make_response(e.message, HTTPStatus.BAD_REQUEST)
@@ -787,13 +808,22 @@ class LayoutObsAPI(Resource):
}
)
def get(self):
content_type = JSON_MIMETYPE
preferred_mimetype = request.accept_mimetypes.best_match(
["application/json", "application/octet-stream"],
"application/json"
)
try:
layout = current_app.data.layout({})
if preferred_mimetype == "application/json":
return make_response(current_app.data.layout({}), HTTPStatus.OK, {"Content-Type": "application/json"})
elif preferred_mimetype == "application/octet-stream":
return make_response(current_app.data.layout_to_fbs_matrix(),
HTTPStatus.OK,
{"Content-Type": "application/octet-stream"})
else:
return make_response(f"Unsupported MIME type '{request.accept_mimetypes}'", HTTPStatus.NOT_ACCEPTABLE)
except PrepareError as e:
return make_response(e.message, HTTPStatus.INTERNAL_SERVER_ERROR)
try:
return make_response(layout, HTTPStatus.OK, {"Content-Type": content_type})
except JSONEncodingValueError as e:
# JSON encoding failure, usually due to bad data
warnings.warn(JSON_NaN_to_num_warning_msg)
+24 -11
View File
@@ -9,18 +9,31 @@ def _mean_var_n(X):
than naive methods (and same method used by numpy.var())
https://en.wikipedia.org/wiki/Algorithms_for_calculating_variance#Two-pass
"""
n = X.shape[0]
if sparse.issparse(X):
mean = X.mean(axis=0).A1
dfm = X - mean
sumsq = np.sum(np.multiply(dfm, dfm), axis=0).A1
v = sumsq / (n - 1)
else:
mean = X.mean(axis=0)
dfm = X - mean
sumsq = np.sum(np.multiply(dfm, dfm), axis=0)
v = sumsq / (n - 1)
# fp_err_occurred is a flag indicating that a floating point error
# occured somewhere in our compute. Used to trigger non-finite
# number handling.
fp_err_occurred = False
def fp_err_set(err, flag):
nonlocal fp_err_occurred
fp_err_occurred = True
with np.errstate(divide="call", invalid="call", call=fp_err_set):
n = X.shape[0]
if sparse.issparse(X):
mean = X.mean(axis=0).A1
dfm = X - mean
sumsq = np.sum(np.multiply(dfm, dfm), axis=0).A1
v = sumsq / (n - 1)
else:
mean = X.mean(axis=0)
dfm = X - mean
sumsq = np.sum(np.multiply(dfm, dfm), axis=0)
v = sumsq / (n - 1)
if fp_err_occurred:
mean[np.isfinite(mean) == False] = 0 # noqa: E712
v[np.isfinite(v) == False] = 0 # noqa: E712
return mean, v, n
+51
View File
@@ -17,6 +17,7 @@ from server.app.util.errors import (
)
from server.app.util.utils import jsonify_scanpy
from server.app.scanpy_engine.diffexp import diffexp_ttest
from server.app.util.fbs.matrix import encode_matrix_fbs
"""
Sort order for methods
@@ -411,6 +412,15 @@ class ScanpyEngine(CXGDriver):
except ValueError:
raise JSONEncodingValueError("Error encoding annotations to JSON")
def annotation_to_fbs_matrix(self, axis, fields=None):
if axis == Axis.OBS:
df = self.data.obs
else:
df = self.data.var
if fields is not None and len(fields) > 0:
df = df[fields]
return encode_matrix_fbs(df, col_idx=df.columns)
def data_frame(self, filter, axis):
"""
Retrieves data for each variable for observations in data frame
@@ -449,6 +459,30 @@ class ScanpyEngine(CXGDriver):
except ValueError:
raise JSONEncodingValueError("Error encoding dataframe to JSON")
def data_frame_to_fbs_matrix(self, filter, axis):
"""
Retrieves data 'X' and returns in a flatbuffer Matrix.
:param filter: filter: dictionary with filter params
:param axis: string obs or var
:return: flatbuffer Matrix
Caveats:
* currently only supports access on VAR axis
* currently only supports filtering on VAR axis
"""
if axis != Axis.VAR:
raise ValueError("Only VAR dimension access is supported")
try:
obs_selector, var_selector = self._filter_to_mask(filter, use_slices=False)
except (KeyError, IndexError) as e:
raise FilterError(f"Error parsing filter: {e}") from e
if obs_selector is not None:
raise FilterError("filtering on obs unsupported")
# Currently only handles VAR dimension
X = self.data._X[:, var_selector]
return encode_matrix_fbs(X, col_idx=np.nonzero(var_selector)[0], row_idx=None)
def diffexp_topN(self, obsFilterA, obsFilterB, top_n=None, interactive_limit=None):
if Axis.VAR in obsFilterA or Axis.VAR in obsFilterB:
raise FilterError("Observation filters may not contain vaiable conditions")
@@ -516,3 +550,20 @@ class ScanpyEngine(CXGDriver):
)
except ValueError:
raise JSONEncodingValueError("Error encoding layout to JSON")
def layout_to_fbs_matrix(self):
"""
Return the default 2-D layout for cells as a FBS Matrix.
Caveats:
* does not support filtering
* only returns Matrix in columnar layout
"""
try:
df_layout = self.data.obsm[f"X_{self.layout_method}"]
except ValueError as e:
raise PrepareError(
f"Layout has not been calculated using {self.layout_method}, "
f"please prepare your datafile and relaunch cellxgene") from e
normalized_layout = (df_layout - df_layout.min()) / (df_layout.max() - df_layout.min())
return encode_matrix_fbs(normalized_layout.astype(dtype=np.float32), col_idx=None, row_idx=None)
-3
View File
@@ -3,9 +3,6 @@ from enum import Enum
DEFAULT_TOP_N = 10
# response mimetypes
JSON_MIMETYPE = "application/json"
class AugmentedEnum(Enum):
def __hash__(self):
+41
View File
@@ -0,0 +1,41 @@
# automatically generated by the FlatBuffers compiler, do not modify
# namespace: NetEncoding
import flatbuffers
class Column(object):
__slots__ = ['_tab']
@classmethod
def GetRootAsColumn(cls, buf, offset):
n = flatbuffers.encode.Get(flatbuffers.packer.uoffset, buf, offset)
x = Column()
x.Init(buf, n + offset)
return x
# Column
def Init(self, buf, pos):
self._tab = flatbuffers.table.Table(buf, pos)
# Column
def UType(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
return self._tab.Get(flatbuffers.number_types.Uint8Flags, o + self._tab.Pos)
return 0
# Column
def U(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(6))
if o != 0:
from flatbuffers.table import Table
obj = Table(bytearray(), 0)
self._tab.Union(obj, o)
return obj
return None
def ColumnStart(builder): builder.StartObject(2)
def ColumnAddUType(builder, uType): builder.PrependUint8Slot(0, uType, 0)
def ColumnAddU(builder, u): builder.PrependUOffsetTRelativeSlot(1, flatbuffers.number_types.UOffsetTFlags.py_type(u), 0)
def ColumnEnd(builder): return builder.EndObject()
@@ -0,0 +1,46 @@
# automatically generated by the FlatBuffers compiler, do not modify
# namespace: NetEncoding
import flatbuffers
class Float32Array(object):
__slots__ = ['_tab']
@classmethod
def GetRootAsFloat32Array(cls, buf, offset):
n = flatbuffers.encode.Get(flatbuffers.packer.uoffset, buf, offset)
x = Float32Array()
x.Init(buf, n + offset)
return x
# Float32Array
def Init(self, buf, pos):
self._tab = flatbuffers.table.Table(buf, pos)
# Float32Array
def Data(self, j):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
a = self._tab.Vector(o)
return self._tab.Get(flatbuffers.number_types.Float32Flags, a + flatbuffers.number_types.UOffsetTFlags.py_type(j * 4))
return 0
# Float32Array
def DataAsNumpy(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
return self._tab.GetVectorAsNumpy(flatbuffers.number_types.Float32Flags, o)
return 0
# Float32Array
def DataLength(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
return self._tab.VectorLen(o)
return 0
def Float32ArrayStart(builder): builder.StartObject(1)
def Float32ArrayAddData(builder, data): builder.PrependUOffsetTRelativeSlot(0, flatbuffers.number_types.UOffsetTFlags.py_type(data), 0)
def Float32ArrayStartDataVector(builder, numElems): return builder.StartVector(4, numElems, 4)
def Float32ArrayEnd(builder): return builder.EndObject()
@@ -0,0 +1,46 @@
# automatically generated by the FlatBuffers compiler, do not modify
# namespace: NetEncoding
import flatbuffers
class Float64Array(object):
__slots__ = ['_tab']
@classmethod
def GetRootAsFloat64Array(cls, buf, offset):
n = flatbuffers.encode.Get(flatbuffers.packer.uoffset, buf, offset)
x = Float64Array()
x.Init(buf, n + offset)
return x
# Float64Array
def Init(self, buf, pos):
self._tab = flatbuffers.table.Table(buf, pos)
# Float64Array
def Data(self, j):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
a = self._tab.Vector(o)
return self._tab.Get(flatbuffers.number_types.Float64Flags, a + flatbuffers.number_types.UOffsetTFlags.py_type(j * 8))
return 0
# Float64Array
def DataAsNumpy(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
return self._tab.GetVectorAsNumpy(flatbuffers.number_types.Float64Flags, o)
return 0
# Float64Array
def DataLength(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
return self._tab.VectorLen(o)
return 0
def Float64ArrayStart(builder): builder.StartObject(1)
def Float64ArrayAddData(builder, data): builder.PrependUOffsetTRelativeSlot(0, flatbuffers.number_types.UOffsetTFlags.py_type(data), 0)
def Float64ArrayStartDataVector(builder, numElems): return builder.StartVector(8, numElems, 8)
def Float64ArrayEnd(builder): return builder.EndObject()
@@ -0,0 +1,46 @@
# automatically generated by the FlatBuffers compiler, do not modify
# namespace: NetEncoding
import flatbuffers
class Int32Array(object):
__slots__ = ['_tab']
@classmethod
def GetRootAsInt32Array(cls, buf, offset):
n = flatbuffers.encode.Get(flatbuffers.packer.uoffset, buf, offset)
x = Int32Array()
x.Init(buf, n + offset)
return x
# Int32Array
def Init(self, buf, pos):
self._tab = flatbuffers.table.Table(buf, pos)
# Int32Array
def Data(self, j):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
a = self._tab.Vector(o)
return self._tab.Get(flatbuffers.number_types.Int32Flags, a + flatbuffers.number_types.UOffsetTFlags.py_type(j * 4))
return 0
# Int32Array
def DataAsNumpy(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
return self._tab.GetVectorAsNumpy(flatbuffers.number_types.Int32Flags, o)
return 0
# Int32Array
def DataLength(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
return self._tab.VectorLen(o)
return 0
def Int32ArrayStart(builder): builder.StartObject(1)
def Int32ArrayAddData(builder, data): builder.PrependUOffsetTRelativeSlot(0, flatbuffers.number_types.UOffsetTFlags.py_type(data), 0)
def Int32ArrayStartDataVector(builder, numElems): return builder.StartVector(4, numElems, 4)
def Int32ArrayEnd(builder): return builder.EndObject()
@@ -0,0 +1,46 @@
# automatically generated by the FlatBuffers compiler, do not modify
# namespace: NetEncoding
import flatbuffers
class JSONEncodedArray(object):
__slots__ = ['_tab']
@classmethod
def GetRootAsJSONEncodedArray(cls, buf, offset):
n = flatbuffers.encode.Get(flatbuffers.packer.uoffset, buf, offset)
x = JSONEncodedArray()
x.Init(buf, n + offset)
return x
# JSONEncodedArray
def Init(self, buf, pos):
self._tab = flatbuffers.table.Table(buf, pos)
# JSONEncodedArray
def Data(self, j):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
a = self._tab.Vector(o)
return self._tab.Get(flatbuffers.number_types.Uint8Flags, a + flatbuffers.number_types.UOffsetTFlags.py_type(j * 1))
return 0
# JSONEncodedArray
def DataAsNumpy(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
return self._tab.GetVectorAsNumpy(flatbuffers.number_types.Uint8Flags, o)
return 0
# JSONEncodedArray
def DataLength(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
return self._tab.VectorLen(o)
return 0
def JSONEncodedArrayStart(builder): builder.StartObject(1)
def JSONEncodedArrayAddData(builder, data): builder.PrependUOffsetTRelativeSlot(0, flatbuffers.number_types.UOffsetTFlags.py_type(data), 0)
def JSONEncodedArrayStartDataVector(builder, numElems): return builder.StartVector(1, numElems, 1)
def JSONEncodedArrayEnd(builder): return builder.EndObject()
+98
View File
@@ -0,0 +1,98 @@
# automatically generated by the FlatBuffers compiler, do not modify
# namespace: NetEncoding
import flatbuffers
class Matrix(object):
__slots__ = ['_tab']
@classmethod
def GetRootAsMatrix(cls, buf, offset):
n = flatbuffers.encode.Get(flatbuffers.packer.uoffset, buf, offset)
x = Matrix()
x.Init(buf, n + offset)
return x
# Matrix
def Init(self, buf, pos):
self._tab = flatbuffers.table.Table(buf, pos)
# Matrix
def NRows(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
return self._tab.Get(flatbuffers.number_types.Uint32Flags, o + self._tab.Pos)
return 0
# Matrix
def NCols(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(6))
if o != 0:
return self._tab.Get(flatbuffers.number_types.Uint32Flags, o + self._tab.Pos)
return 0
# Matrix
def Columns(self, j):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(8))
if o != 0:
x = self._tab.Vector(o)
x += flatbuffers.number_types.UOffsetTFlags.py_type(j) * 4
x = self._tab.Indirect(x)
from .Column import Column
obj = Column()
obj.Init(self._tab.Bytes, x)
return obj
return None
# Matrix
def ColumnsLength(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(8))
if o != 0:
return self._tab.VectorLen(o)
return 0
# Matrix
def ColIndexType(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(10))
if o != 0:
return self._tab.Get(flatbuffers.number_types.Uint8Flags, o + self._tab.Pos)
return 0
# Matrix
def ColIndex(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(12))
if o != 0:
from flatbuffers.table import Table
obj = Table(bytearray(), 0)
self._tab.Union(obj, o)
return obj
return None
# Matrix
def RowIndexType(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(14))
if o != 0:
return self._tab.Get(flatbuffers.number_types.Uint8Flags, o + self._tab.Pos)
return 0
# Matrix
def RowIndex(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(16))
if o != 0:
from flatbuffers.table import Table
obj = Table(bytearray(), 0)
self._tab.Union(obj, o)
return obj
return None
def MatrixStart(builder): builder.StartObject(7)
def MatrixAddNRows(builder, nRows): builder.PrependUint32Slot(0, nRows, 0)
def MatrixAddNCols(builder, nCols): builder.PrependUint32Slot(1, nCols, 0)
def MatrixAddColumns(builder, columns): builder.PrependUOffsetTRelativeSlot(2, flatbuffers.number_types.UOffsetTFlags.py_type(columns), 0)
def MatrixStartColumnsVector(builder, numElems): return builder.StartVector(4, numElems, 4)
def MatrixAddColIndexType(builder, colIndexType): builder.PrependUint8Slot(3, colIndexType, 0)
def MatrixAddColIndex(builder, colIndex): builder.PrependUOffsetTRelativeSlot(4, flatbuffers.number_types.UOffsetTFlags.py_type(colIndex), 0)
def MatrixAddRowIndexType(builder, rowIndexType): builder.PrependUint8Slot(5, rowIndexType, 0)
def MatrixAddRowIndex(builder, rowIndex): builder.PrependUOffsetTRelativeSlot(6, flatbuffers.number_types.UOffsetTFlags.py_type(rowIndex), 0)
def MatrixEnd(builder): return builder.EndObject()
@@ -0,0 +1,12 @@
# automatically generated by the FlatBuffers compiler, do not modify
# namespace: NetEncoding
class TypedArray(object):
NONE = 0
Float32Array = 1
Int32Array = 2
Uint32Array = 3
Float64Array = 4
JSONEncodedArray = 5
@@ -0,0 +1,46 @@
# automatically generated by the FlatBuffers compiler, do not modify
# namespace: NetEncoding
import flatbuffers
class Uint32Array(object):
__slots__ = ['_tab']
@classmethod
def GetRootAsUint32Array(cls, buf, offset):
n = flatbuffers.encode.Get(flatbuffers.packer.uoffset, buf, offset)
x = Uint32Array()
x.Init(buf, n + offset)
return x
# Uint32Array
def Init(self, buf, pos):
self._tab = flatbuffers.table.Table(buf, pos)
# Uint32Array
def Data(self, j):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
a = self._tab.Vector(o)
return self._tab.Get(flatbuffers.number_types.Uint32Flags, a + flatbuffers.number_types.UOffsetTFlags.py_type(j * 4))
return 0
# Uint32Array
def DataAsNumpy(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
return self._tab.GetVectorAsNumpy(flatbuffers.number_types.Uint32Flags, o)
return 0
# Uint32Array
def DataLength(self):
o = flatbuffers.number_types.UOffsetTFlags.py_type(self._tab.Offset(4))
if o != 0:
return self._tab.VectorLen(o)
return 0
def Uint32ArrayStart(builder): builder.StartObject(1)
def Uint32ArrayAddData(builder, data): builder.PrependUOffsetTRelativeSlot(0, flatbuffers.number_types.UOffsetTFlags.py_type(data), 0)
def Uint32ArrayStartDataVector(builder, numElems): return builder.StartVector(4, numElems, 4)
def Uint32ArrayEnd(builder): return builder.EndObject()
+207
View File
@@ -0,0 +1,207 @@
import flatbuffers
import numpy as np
from scipy import sparse
import pandas as pd
import server.app.util.fbs.NetEncoding.Column as Column
import server.app.util.fbs.NetEncoding.TypedArray as TypedArray
import server.app.util.fbs.NetEncoding.Matrix as Matrix
# Placeholder until recent enhancements to flatbuffers Python
# runtime are released, at which point we can use the default
# version. This code is a port of the head. See:
#
# https://github.com/google/flatbuffers/pull/4829
#
def CreateNumpyVector(builder, x):
"""CreateNumpyVector writes a numpy array into the buffer."""
if not isinstance(x, np.ndarray):
raise TypeError("non-numpy-ndarray passed to CreateNumpyVector")
if x.dtype.kind not in ['b', 'i', 'u', 'f']:
raise TypeError("numpy-ndarray holds elements of unsupported datatype")
if x.ndim > 1:
raise TypeError("multidimensional-ndarray passed to CreateNumpyVector")
builder.StartVector(x.itemsize, x.size, x.dtype.alignment)
# Ensure little endian byte ordering
if x.dtype.str[0] == "<":
x_little_endian = x
else:
x_little_endian = x.byteswap(inplace=False)
# Calculate total length
len = int(x_little_endian.itemsize * x_little_endian.size)
builder.head = int(builder.Head() - len)
# tobytes ensures c_contiguous ordering
builder.Bytes[builder.Head():builder.Head() + len] = x_little_endian.tobytes(order='C')
return builder.EndVector(x.size)
# Serialization helper
def serialize_column(builder, typed_arr):
""" Serialize NetEncoding.Column """
(u_type, u_value) = typed_arr
Column.ColumnStart(builder)
Column.ColumnAddUType(builder, u_type)
Column.ColumnAddU(builder, u_value)
return Column.ColumnEnd(builder)
# Serialization helper
def serialize_matrix(builder, n_rows, n_cols, columns, col_idx):
""" Serialize NetEncoding.Matrix """
Matrix.MatrixStart(builder)
Matrix.MatrixAddNRows(builder, n_rows)
Matrix.MatrixAddNCols(builder, n_cols)
Matrix.MatrixAddColumns(builder, columns)
if col_idx is not None:
(u_type, u_val) = col_idx
Matrix.MatrixAddColIndexType(builder, u_type)
Matrix.MatrixAddColIndex(builder, u_val)
return Matrix.MatrixEnd(builder)
# Serialization helper
def serialize_typed_array(builder, source_array, encoding_info):
"""
Serialize any of the various typed arrays, eg, Float32Array. Specific
means of serialization and type conversion are provided by type_info.
"""
arr = source_array
(array_type, as_type) = encoding_info(source_array)
if isinstance(arr, pd.Index):
arr = arr.to_series()
# convert to a simple ndarray
if as_type == 'json':
as_json = arr.to_json(orient='records')
arr = np.array(bytearray(as_json, 'utf-8'))
else:
if sparse.issparse(arr):
arr = arr.toarray()
elif isinstance(arr, pd.Series):
arr = arr.get_values()
if arr.dtype != as_type:
arr = arr.astype(as_type)
# serialize the ndarray into a vector
if arr.ndim == 2 and arr.shape[0] == 1:
arr = arr[0]
vec = CreateNumpyVector(builder, arr)
# serialize the typed array table
builder.StartObject(1)
builder.PrependUOffsetTRelativeSlot(0, vec, 0)
array_value = builder.EndObject()
return (array_type, array_value)
def column_encoding(arr):
type_map = {
# dtype: ( array_type, as_type )
np.float64: (TypedArray.TypedArray.Float32Array, np.float32),
np.float32: (TypedArray.TypedArray.Float32Array, np.float32),
np.float16: (TypedArray.TypedArray.Float32Array, np.float32),
np.int8: (TypedArray.TypedArray.Int32Array, np.int32),
np.int16: (TypedArray.TypedArray.Int32Array, np.int32),
np.int32: (TypedArray.TypedArray.Int32Array, np.int32),
np.int64: (TypedArray.TypedArray.Int32Array, np.int32),
np.uint8: (TypedArray.TypedArray.Uint32Array, np.uint32),
np.uint16: (TypedArray.TypedArray.Uint32Array, np.uint32),
np.uint32: (TypedArray.TypedArray.Uint32Array, np.uint32),
np.uint64: (TypedArray.TypedArray.Uint32Array, np.uint32)
}
type_map_default = (TypedArray.TypedArray.JSONEncodedArray, 'json')
return type_map.get(arr.dtype.type, type_map_default)
def index_encoding(arr):
type_map = {
# dtype: ( array_type, as_type )
np.int32: (TypedArray.TypedArray.Int32Array, np.int32),
np.int64: (TypedArray.TypedArray.Int32Array, np.int32),
np.uint32: (TypedArray.TypedArray.Uint32Array, np.uint32),
np.uint64: (TypedArray.TypedArray.Uint32Array, np.uint32)
}
type_map_default = (TypedArray.TypedArray.JSONEncodedArray, 'json')
return type_map.get(arr.dtype.type, type_map_default)
def guess_at_mem_needed(matrix):
(n_rows, n_cols) = matrix.shape
if isinstance(matrix, np.ndarray) or sparse.issparse(matrix):
guess = (n_rows * n_cols * matrix.dtype.itemsize) + 1024
elif isinstance(matrix, pd.DataFrame):
# XXX TODO - DataFrame type estimate
guess = 1
else:
guess = 1
# round up to nearest 1024 bytes
guess = (guess + 0x400) & (~0x3ff)
return guess
def encode_matrix_fbs(matrix, row_idx=None, col_idx=None):
"""
Given a 2D DataFrame, ndarray or sparse equivalent, create and return a
Matrix flatbuffer.
:param matrix: 2D DataFrame, ndarray or sparse equivalent
:param row_idx: index for row dimension, Index or ndarray
:param col_idx: index for col dimension, Index or ndarray
NOTE: row indices are (currently) unsupported and must be None
"""
if row_idx is not None:
raise ValueError("row indexing not supported for FBS Matrix")
if matrix.ndim != 2:
raise ValueError("FBS Matrix must be 2D")
(n_rows, n_cols) = matrix.shape
# estimate size needed, so we don't unnecessarily realloc.
builder = flatbuffers.Builder(guess_at_mem_needed(matrix))
if isinstance(matrix, pd.DataFrame):
matrix_columns = reversed(tuple(matrix[name] for name in matrix))
else:
matrix_columns = reversed(tuple(c for c in matrix.T))
columns = []
# for idx in reversed(np.arange(n_cols)):
for c in matrix_columns:
# serialize the typed array
typed_arr = serialize_typed_array(builder, c, column_encoding)
# serialize the Column union
columns.append(serialize_column(builder, typed_arr))
# Serialize Matrix.columns[]
Matrix.MatrixStartColumnsVector(builder, n_cols)
for c in columns:
builder.PrependUOffsetTRelative(c)
matrix_column_vec = builder.EndVector(n_cols)
# serialize the colIndex if provided
cidx = None
if col_idx is not None:
cidx = serialize_typed_array(builder, col_idx, index_encoding)
# Serialize Matrix
matrix = serialize_matrix(builder, n_rows, n_cols, matrix_column_vec, cidx)
builder.Finish(matrix)
return builder.Output()