mirror of
https://github.com/paulscherrerinstitute/sf_daq_broker.git
synced 2026-09-02 09:10:43 +02:00
retrieve from data buffer using data_api3, currently writing BSDATA file in parallel to BSREAD
This commit is contained in:
committed by
Data Backend account
parent
d39a411507
commit
fe70e61aa5
@@ -190,6 +190,10 @@ class BrokerManager(object):
|
||||
request.get("channels_list"),
|
||||
config.OUTPUT_FILE_SUFFIX_DATA_BUFFER)
|
||||
|
||||
send_write_request(broker_config.TAG_DATA3BUFFER,
|
||||
request.get("channels_list"),
|
||||
config.OUTPUT_FILE_SUFFIX_DATA3_BUFFER)
|
||||
|
||||
send_write_request(broker_config.TAG_IMAGEBUFFER,
|
||||
request.get("camera_list"),
|
||||
config.OUTPUT_FILE_SUFFIX_IMAGE_BUFFER)
|
||||
|
||||
@@ -10,9 +10,10 @@ AUDIT_FILE_TIME_FORMAT = "%Y%m%d-%H%M%S"
|
||||
|
||||
DEFAULT_AUDIT_FILENAME = "/var/log/sf_databuffer_audit.log"
|
||||
|
||||
DATA_API_QUERY_ADDRESS = "http://sf-data-api-02.psi.ch/query"
|
||||
#DATA_API_QUERY_ADDRESS = "http://sf-data-api-02.psi.ch/query"
|
||||
DATA_API_QUERY_ADDRESS = "https://data-api.psi.ch/sf/query"
|
||||
IMAGE_API_QUERY_ADDRESS = ["http://172.27.0.14:8080/api/v1/query", "http://172.27.0.15:8080/api/v1/query"]
|
||||
DATA_API3_QUERY_ADDRESS = "http://sf-daqbuf-21:8371/api/1.0.1/query"
|
||||
DATA_API3_QUERY_ADDRESS = "http://sf-daqbuf-33:8371/api/1"
|
||||
EPICS_QUERY_ADDRESS = "https://data-api.psi.ch/sf"
|
||||
|
||||
DATA_BACKEND = "sf-databuffer"
|
||||
@@ -23,5 +24,6 @@ TRANSFORM_PULSE_ID_TO_TIMESTAMP_QUERY = False
|
||||
SEPARATE_CAMERA_CHANNELS = True
|
||||
|
||||
OUTPUT_FILE_SUFFIX_DATA_BUFFER = "BSREAD"
|
||||
OUTPUT_FILE_SUFFIX_DATA3_BUFFER = "BSDATA"
|
||||
OUTPUT_FILE_SUFFIX_IMAGE_BUFFER = "CAMERAS"
|
||||
OUTPUT_FILE_SUFFIX_EPICS = "PVCHANNELS"
|
||||
|
||||
@@ -11,4 +11,5 @@ DEFAULT_QUEUE = "write_request"
|
||||
# Name of the queue for bsread write requests.
|
||||
TAG_DATABUFFER = "databuffer"
|
||||
TAG_IMAGEBUFFER = "imagebuffer"
|
||||
TAG_EPICS = "epics"
|
||||
TAG_EPICS = "epics"
|
||||
TAG_DATA3BUFFER = "data3buffer"
|
||||
|
||||
@@ -6,7 +6,7 @@ import requests
|
||||
|
||||
from sf_daq_broker import config
|
||||
|
||||
_logger = getLogger(__name__)
|
||||
_logger = getLogger("broker_writer")
|
||||
|
||||
|
||||
def get_data_api_request(channels, start_pulse_id, stop_pulse_id):
|
||||
@@ -50,7 +50,7 @@ def transform_range_from_pulse_id_to_timestamp(data_api_request):
|
||||
mapping_request = {'range': {'startPulseId': data_api_request["range"]["startPulseId"],
|
||||
'endPulseId': data_api_request["range"]["endPulseId"]+1}}
|
||||
|
||||
mapping_response = requests.post(url=config.DATA_API_QUERY_ADDRESS + "/mapping", json=mapping_request).json()
|
||||
mapping_response = requests.post(url=config.DATA_API_QUERY_ADDRESS + "/mapping", json=mapping_request, timeout=10).json()
|
||||
|
||||
_logger.info("Response to mapping request: %s", mapping_response)
|
||||
|
||||
@@ -63,6 +63,7 @@ def transform_range_from_pulse_id_to_timestamp(data_api_request):
|
||||
_logger.info("Transformed request to startSeconds and endSeconds. %s" % new_data_api_request)
|
||||
|
||||
except Exception as e:
|
||||
_logger.error(e)
|
||||
raise RuntimeError("Cannot retrieve the pulse_id to timestamp mapping.") from e
|
||||
|
||||
return new_data_api_request
|
||||
|
||||
@@ -7,6 +7,7 @@ import h5py
|
||||
import numpy
|
||||
import requests
|
||||
from random import randrange
|
||||
from copy import deepcopy
|
||||
|
||||
from sf_daq_broker import config, utils
|
||||
from sf_daq_broker.writer.utils import channel_type_deserializer_mapping
|
||||
@@ -29,7 +30,11 @@ def write_from_databuffer(data_api_request, output_file, metadata):
|
||||
|
||||
start_time = time()
|
||||
|
||||
response = requests.post(url=config.DATA_API_QUERY_ADDRESS, json=data_api_request)
|
||||
new_data_api_request = deepcopy(data_api_request)
|
||||
# new_data_api_request["range"]["startPulseId"] -= 1
|
||||
# new_data_api_request["range"]["endPulseId"] += 1
|
||||
|
||||
response = requests.post(url=config.DATA_API_QUERY_ADDRESS, json=new_data_api_request, timeout=1000)
|
||||
data = json.loads(response.content)
|
||||
|
||||
if not data:
|
||||
@@ -91,6 +96,45 @@ def write_from_imagebuffer(data_api_request, output_file, parameters):
|
||||
_logger.error(e)
|
||||
|
||||
|
||||
def write_from_databuffer_api3(data_api_request, output_file, parameters):
|
||||
import data_api3.h5 as h5
|
||||
import pytz
|
||||
|
||||
_logger.debug("Data3 API request: %s", data_api_request)
|
||||
|
||||
data_api_request_timestamp = utils.transform_range_from_pulse_id_to_timestamp(data_api_request)
|
||||
|
||||
channels = [channel["name"] for channel in data_api_request_timestamp["channels"]]
|
||||
|
||||
start = datetime.fromtimestamp(float(data_api_request_timestamp["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_timestamp["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
|
||||
}
|
||||
}
|
||||
|
||||
data_buffer_url = config.DATA_API3_QUERY_ADDRESS
|
||||
|
||||
_logger.debug("Requesting '%s' to output_file %s from %s " %
|
||||
(query, output_file, data_buffer_url))
|
||||
|
||||
start_time = time()
|
||||
|
||||
try:
|
||||
h5.request(query, filename=output_file, baseurl=data_buffer_url, default_backend=config.DATA_BACKEND)
|
||||
_logger.info("Data download and writing took %s seconds." % (time() - start_time))
|
||||
except Exception as e:
|
||||
_logger.error("Got exception from data_api3")
|
||||
_logger.error(e)
|
||||
|
||||
|
||||
class BsreadH5Writer(object):
|
||||
|
||||
def __init__(self, output_file, metadata):
|
||||
|
||||
@@ -11,7 +11,7 @@ from pika import BlockingConnection, ConnectionParameters, BasicProperties
|
||||
from sf_daq_broker import config, utils
|
||||
import sf_daq_broker.rabbitmq.config as broker_config
|
||||
from sf_daq_broker.utils import get_data_api_request
|
||||
from sf_daq_broker.writer.bsread_writer import write_from_imagebuffer, write_from_databuffer
|
||||
from sf_daq_broker.writer.bsread_writer import write_from_imagebuffer, write_from_databuffer, write_from_databuffer_api3
|
||||
from sf_daq_broker.writer.epics_writer import write_epics_pvs
|
||||
|
||||
_logger = logging.getLogger("broker_writer")
|
||||
@@ -70,6 +70,8 @@ def process_request(request):
|
||||
|
||||
if writer_type == broker_config.TAG_DATABUFFER:
|
||||
logger_data_api = logging.getLogger("data_api")
|
||||
elif writer_type == broker_config.TAG_DATA3BUFFER:
|
||||
logger_data_api = logging.getLogger("data_api3")
|
||||
elif writer_type == broker_config.TAG_IMAGEBUFFER:
|
||||
logger_data_api = logging.getLogger("data_api3")
|
||||
elif writer_type == broker_config.TAG_EPICS:
|
||||
@@ -99,6 +101,10 @@ def process_request(request):
|
||||
_logger.info("Using databuffer writer.")
|
||||
write_from_databuffer(get_data_api_request(channels, start_pulse_id, stop_pulse_id), output_file, metadata)
|
||||
|
||||
elif writer_type == broker_config.TAG_DATA3BUFFER:
|
||||
_logger.info("Using data_api3 databuffer writer.")
|
||||
write_from_databuffer_api3(get_data_api_request(channels, start_pulse_id, stop_pulse_id), output_file, metadata)
|
||||
|
||||
elif writer_type == broker_config.TAG_IMAGEBUFFER:
|
||||
_logger.info("Using imagebuffer writer.")
|
||||
write_from_imagebuffer(get_data_api_request(channels, start_pulse_id, stop_pulse_id), output_file, metadata)
|
||||
|
||||
Reference in New Issue
Block a user