Create a synchronization example for TDMS logging #820
Description
Activity
Hmm, I don't think this is possible with our TDMS APIs. But I'm not 100% sure. I'll ponder.
Ah... Thanks for having a ponder. Would it then be best to log to a different TDMS for each task and post-process into one? Had a go but not sure how stable this is, could the callback cause a delay that means the buffer overflows...
with nidaqmx.Task() as ai_task, nidaqmx.Task() as ci1_task, nidaqmx.Task() as ci2_task, nidaqmx.Task() as clk_task: def callback(task_handle, every_n_samples_event_type, number_of_samples, callback_data): """Callback function for reading signals.""" nonlocal total_ai_read nonlocal total_ci1_read nonlocal total_ci2_read ai_read = ai_task.read(number_of_samples_per_channel=number_of_samples) ci1_read = ci1_task.read(number_of_samples_per_channel=number_of_samples) ci2_read = ci1_task.read(number_of_samples_per_channel=number_of_samples) total_ai_read += len(ai_read) total_ci1_read += len(ci1_read) total_ci2_read += len(ci2_read) print(f"\t{len(ai_read)}\t{len(ci1_read)}\t{len(ci2_read)}\t\t{total_ai_read}\t{total_ci1_read}\t{total_ci2_read}", end="\r") return 0 #Configure sample clock clk_task.co_channels.add_co_pulse_chan_freq( counter="Dev1/ctr3", freq=sample_rate, ) clk_task.timing.cfg_implicit_timing(sample_mode=AcquisitionType.CONTINUOUS) #Configure ai channels ai_task.ai_channels.add_ai_voltage_chan("Dev1/ai0", "Torque01") ai_task.ai_channels.add_ai_voltage_chan("Dev1/ai1", "Torque02") ai_task.timing.cfg_samp_clk_timing(sample_rate, "/Dev1/Ctr3InternalOutput", sample_mode=AcquisitionType.CONTINUOUS) ai_task.register_every_n_samples_acquired_into_buffer_event(sample_rate, callback) # Configure ci channels #ci 1 ci1_chan = ci1_task.ci_channels.add_ci_count_edges_chan( "Dev1/ctr0", "Rotations01", edge=Edge.RISING, initial_count=0 ) ci1_chan.ci_count_edges_term = "/Dev1/PFI0" ci1_task.timing.cfg_samp_clk_timing( sample_rate, "/Dev1/Ctr3InternalOutput", sample_mode=AcquisitionType.CONTINUOUS ) #ci 2 ci2_chan = ci2_task.ci_channels.add_ci_count_edges_chan( "Dev1/ctr1", "Rotations02", edge=Edge.RISING, initial_count=0 ) ci2_chan.ci_count_edges_term = "/Dev1/PFI1" ci2_task.timing.cfg_samp_clk_timing( sample_rate, "/Dev1/Ctr3InternalOutput", sample_mode=AcquisitionType.CONTINUOUS ) #configure logging ai_task.in_stream.configure_logging( "{0}_ai.tdms".format(filepath), LoggingMode.LOG_AND_READ, operation=LoggingOperation.CREATE_OR_REPLACE ) ci1_task.in_stream.configure_logging( "{0}_ci1.tdms".format(filepath), LoggingMode.LOG_AND_READ, operation=LoggingOperation.CREATE_OR_REPLACE ) ci2_task.in_stream.configure_logging( "{0}_ci2.tdms".format(filepath), LoggingMode.LOG_AND_READ, operation=LoggingOperation.CREATE_OR_REPLACE ) clk_task.start() ci1_task.start() ci2_task.start() ai_task.start() print("Acquiring samples continuously. Press Enter to stop.\n") print("Read:\tAI\tCI1\tCI2\tTotal:\tAI\tCI1\tCI2") input() ai_task.stop() ci1_task.stop() ci2_task.stop() clk_task.stop() print(f"\nAcquired {total_ai_read} total AI samples and {total_ci1_read} total CI1 samples.")
That approach looks sound. Your example isn't combining, but you'll end up with 3 TDMS files that you could merge as a post-process. Some thoughts:
- You're synchronizing by using a shared CO task - good!
- Yes, if your callbacks are slow you could get behind in your buffer and eventually overflow. But you're not doing any significant processing, so I think you'll be OK. Especially true since you're only getting a callback once per second. I recommend 10x per second or slower.
Logging multiple DAQmx tasks in the same TDMS file can't be done by DAQmx Configure Logging VI. (Reference )
The reference leads to an example VI that shows how to log data to a single TDMS file when coming from more than one DAQmx acquisition tasks, using a producer-consumer queue. (Example VI)
I have written my implementation based on that VI. Would this be sufficient? @zhindes
# Configuration SAMPLE_RATE = 1000 SAMPLES_PER_CHANNEL = 1000 TIMEOUT = 10.0 def producer( tasks: List[nidaqmx.Task], data_queue: queue.Queue, stop_event: threading.Event ) -> None: """Producer function that reads data from DAQmx tasks and puts it in the queue.""" try: while not stop_event.is_set(): # Read from all tasks data = [] for task in tasks: task_data = task.read( number_of_samples_per_channel=SAMPLES_PER_CHANNEL, timeout=TIMEOUT ) data.append(task_data) # Put data in queue data_queue.put(data) except Exception as e: print(f"Error in producer: {e}") stop_event.set() finally: # Signal consumer that we're done data_queue.put(None) def consumer( data_queue: queue.Queue, tdms_path: str, group_names: List[str], channel_names: List[List[str]], stop_event: threading.Event ) -> None: """Consumer function that writes data from the queue to a TDMS file.""" try: with TdmsWriter(tdms_path) as tdms_writer: while not stop_event.is_set(): try: # Get data from queue with timeout data = data_queue.get(timeout=TIMEOUT) # Check for producer completion if data is None: break # Create TDMS objects for each channel root_object = RootObject(properties={ "Creation Time": time.strftime("%Y-%m-%d %H:%M:%S") }) objects_to_write = [root_object] # Write data for each task/group for task_idx, task_data in enumerate(data): group = GroupObject( group_names[task_idx], properties={"Sample Rate": SAMPLE_RATE} ) objects_to_write.append(group) # Convert data to numpy arrays and ensure 1D if isinstance(task_data, (list, tuple)) and isinstance(task_data[0], (list, tuple, np.ndarray)): # Multiple channels (AI task) for chan_idx, chan_data in enumerate(task_data): chan_data = np.array(chan_data).flatten() # Ensure 1D array channel = ChannelObject( group_names[task_idx], channel_names[task_idx][chan_idx], chan_data, properties={"Sample Rate": SAMPLE_RATE} ) objects_to_write.append(channel) else: # Single channel (CI task) task_data = np.array(task_data).flatten() # Ensure 1D array channel = ChannelObject( group_names[task_idx], channel_names[task_idx][0], task_data, properties={"Sample Rate": SAMPLE_RATE} ) objects_to_write.append(channel) # Write to TDMS file tdms_writer.write_segment(objects_to_write) except queue.Empty: continue except Exception as e: print(f"Error in consumer: {e}") stop_event.set() def main(): # Create a queue for data transfer data_queue = queue.Queue(maxsize=10) stop_event = threading.Event() # Create tasks ai_task = nidaqmx.Task() ci1_task = nidaqmx.Task() clk_task = nidaqmx.Task() try: # Configure sample clock clk_task.co_channels.add_co_pulse_chan_freq( counter="Dev3/ctr1", freq=SAMPLE_RATE, ) clk_task.timing.cfg_implicit_timing(sample_mode=AcquisitionType.CONTINUOUS) # Configure AI task ai_task.ai_channels.add_ai_voltage_chan("Dev2/ai0", "Torque01") ai_task.ai_channels.add_ai_voltage_chan("Dev2/ai1", "Torque02") ai_task.timing.cfg_samp_clk_timing( SAMPLE_RATE, sample_mode=AcquisitionType.CONTINUOUS, samps_per_chan=SAMPLES_PER_CHANNEL ) # Configure CI task ci1_chan = ci1_task.ci_channels.add_ci_count_edges_chan( "Dev3/ctr0", "Rotations01", edge=Edge.RISING, initial_count=0 ) ci1_chan.ci_count_edges_term = "/Dev3/PFI0" ci1_task.timing.cfg_samp_clk_timing( SAMPLE_RATE, "/Dev3/Ctr3InternalOutput", sample_mode=AcquisitionType.CONTINUOUS, samps_per_chan=SAMPLES_PER_CHANNEL ) # Create threads producer_thread = threading.Thread( target=producer, args=([ai_task, ci1_task], data_queue, stop_event) ) consumer_thread = threading.Thread( target=consumer, args=( data_queue, "multi_task_data.tdms", ["AI_Task", "CI_Task"], [["Torque01", "Torque02"], ["Rotations01"]], stop_event ) ) # Start tasks in correct order clk_task.start() ci1_task.start() ai_task.start() # Start threads producer_thread.start() consumer_thread.start() print("Acquiring and logging data. Press Enter to stop...") input() # Stop acquisition stop_event.set() # Wait for threads to complete producer_thread.join() consumer_thread.join() finally: # Cleanup for task in [ai_task, ci1_task, clk_task]: if task: task.stop() task.close() print("\nAcquisition complete. Data saved to multi_task_data.tdms")
The examples in synchronisation (required for multiple task types eg. analogue in and counter in) do not show how to log all tasks to a single TDMS file. It would be useful to have all synchronised tasks in one place as the NI hardware supports multiple tasks.
Doing something like this raises the error "File specified is already opened for output. NI-DAQmx requires exclusive write access.
ai_task.in_stream.configure_logging( filepath, LoggingMode.LOG, operation=LoggingOperation.CREATE_OR_REPLACE ) ci1_task.in_stream.configure_logging( filepath, LoggingMode.LOG, operation=LoggingOperation.OPEN )AB#3252273