From e827cf60b84d7fa89869c1430186387c6bdf083a Mon Sep 17 00:00:00 2001 From: Andrej Babic Date: Thu, 20 Aug 2020 11:25:49 +0200 Subject: [PATCH] Move writing methods to bsread_writer --- sf_daq_broker/writer/bsread_writer.py | 72 +++++++++++++++++++++++++++ 1 file changed, 72 insertions(+) diff --git a/sf_daq_broker/writer/bsread_writer.py b/sf_daq_broker/writer/bsread_writer.py index 016a638..938b10e 100644 --- a/sf_daq_broker/writer/bsread_writer.py +++ b/sf_daq_broker/writer/bsread_writer.py @@ -219,3 +219,75 @@ class CompactBsreadH5Writer(BsreadH5Writer): self.file["/data/" + name + "/global_date"] = data["global_date"] self.file["/data/" + name + "/data"] = data["data"] self.file["/data/" + name + "/is_data_present"] = data["is_data_present"] + + +def write_from_databuffer(data_api_request, parameters): + data, data_len = get_data_from_buffer(data_api_request) + _logger.info("Data retrieval (%d bytes) took %s seconds." % (data_len, time() - start_time)) + + start_time = time() + write_data_to_file(parameters, data) + _logger.info("Data writing took %s seconds." % (time() - start_time)) + + +def write_from_imagebuffer(data_api_request_pulseid, parameters): + import data_api3.h5 as h5 + import pytz + + data_api_request = utils.transform_range_from_pulse_id_to_timestamp(data_api_request_pulseid) + + channels = [channel["name"] for channel in data_api_request["channels"]] + + filename = parameters["output_file"] + + start = datetime.fromtimestamp(float(data_api_request["range"]["startSeconds"])).astimezone(pytz.timezone('UTC')).strftime("%Y-%m-%dT%H:%M:%S.%fZ") # isoformat() # "2019-12-13T09:00:00.00 + end = datetime.fromtimestamp(float(data_api_request["range"]["endSeconds"])).astimezone(pytz.timezone('UTC')).strftime("%Y-%m-%dT%H:%M:%S.%fZ") # isoformat() # "2019-12-13T09:00:00.00 + + query = { + "channels": channels, + "range": { + "type": "date", + "startDate": start, + "endDate": end + } + } + + _logger.info("Going to make query %s to write file %s from %s " % (query, filename, config.IMAGE_API_QUERY_ADDRESS)) + + h5.request(query, filename, url=config.IMAGE_API_QUERY_ADDRESS) + + +def write_data_to_file(parameters, json_data): + + if not parameters: + raise ValueError("Received parameters from broker are empty. parameters=%s" % parameters) + + output_file = parameters["output_file"] + + writer = CompactBsreadH5Writer(output_file, parameters) + + writer.write_data(json_data) + writer.close() + + +def get_data_from_buffer(data_api_request): + + _logger.info("Loading data for range: %s" % data_api_request["range"]) + + _logger.debug("Data API request: %s", data_api_request) + + response = requests.post(url=config.DATA_API_QUERY_ADDRESS, json=data_api_request) + + data, data_len = json.loads(response.content), len(response.content) + + if not data: + raise ValueError("Received data from data_api is empty. data=%s" % data) + + # The response is a list if status is OK, otherwise its a dictionary, of course. + if isinstance(data, dict): + if data.get("status") == 500: + raise Exception("Server returned error: %s" % data) + + raise Exception("Server returned a dict (keys: %s), but a list was expected." % list(data.keys())) + + return data, data_len