Skip to main content

Buffered FIFO queue for streaming bytes in python?

I am doing a lot of processing and generating of csv and json files in python, to which I'm new-ish. Here's what I have so far:

row_iter = await get_row_iterator() # basically returns a db cursor for streaming rows
gcs_object = self.gcs_client.bucket(self.env.temp_bucket).blob(manifest_key)

with gcs_object.open(mode="w", content_type="text/csv") as f:
    async for row_obj in row_iter:
        f.write(f"{row_obj.some_column},{row_obj.some_other_column}\n")

This works fine, but if I understand correctly, the problem is that there is a lot of waiting around. While I'm writing a row with f.write(), the row-reading iterator row_iter is just sitting there doing nothing. While I'm reading from the iterator, the writer is twiddling its thumbs.

While the way I've done this conserves memory, from a speed perspective its no better than reading all the rows into memory first, and then starting an upload. There should be no reason I can't be uploading data that has already been processed while also downloading and processing the next batch. I'd like to stream this process, so that the row iterator can be reading rows and stuffing them into a fixed size FIFO buffer as fast as it possibly can, while on the other end I'm reading from the buffer and writing the csv to GCS as fast as the network can handle. If the buffer fills up, writing to it should block until there is space. If the buffer empties, reading from it should block until there is enough to read so that the maximum amount of memory used by this pipeline is controlled.

It looks an awful lot like I should be able to do something like this with python's Buffered Streams, but I can't figure out how to use them to accomplish this, and I can't find any examples of how they should be used. I'll need to do this with csv, json, and probably other custom formats, so I'm not necessarily looking for a package that handles this for a specific format, but a more generic way, even a way to write my own streaming pipelines.



source https://stackoverflow.com/questions/75415211/buffered-fifo-queue-for-streaming-bytes-in-python

Comments

Popular posts from this blog

ValueError: X has 10 features, but LinearRegression is expecting 1 features as input

So, I am trying to predict the model but its throwing error like it has 10 features but it expacts only 1. So I am confused can anyone help me with it? more importantly its not working for me when my friend runs it. It works perfectly fine dose anyone know the reason about it? cv = KFold(n_splits = 10) all_loss = [] for i in range(9): # 1st for loop over polynomial orders poly_order = i X_train = make_polynomial(x, poly_order) loss_at_order = [] # initiate a set to collect loss for CV for train_index, test_index in cv.split(X_train): print('TRAIN:', train_index, 'TEST:', test_index) X_train_cv, X_test_cv = X_train[train_index], X_test[test_index] t_train_cv, t_test_cv = t[train_index], t[test_index] reg.fit(X_train_cv, t_train_cv) loss_at_order.append(np.mean((t_test_cv - reg.predict(X_test_cv))**2)) # collect loss at fold all_loss.append(np.mean(loss_at_order)) # collect loss at order plt.plot(np.log(al...

Sorting large arrays of big numeric stings

I was solving bigSorting() problem from hackerrank: Consider an array of numeric strings where each string is a positive number with anywhere from to digits. Sort the array's elements in non-decreasing, or ascending order of their integer values and return the sorted array. I know it works as follows: def bigSorting(unsorted): return sorted(unsorted, key=int) But I didnt guess this approach earlier. Initially I tried below: def bigSorting(unsorted): int_unsorted = [int(i) for i in unsorted] int_sorted = sorted(int_unsorted) return [str(i) for i in int_sorted] However, for some of the test cases, it was showing time limit exceeded. Why is it so? PS: I dont know exactly what those test cases were as hacker rank does not reveal all test cases. source https://stackoverflow.com/questions/73007397/sorting-large-arrays-of-big-numeric-stings

How to load Javascript with imported modules?

I am trying to import modules from tensorflowjs, and below is my code. test.html <!DOCTYPE html> <html lang="en"> <head> <meta charset="UTF-8"> <title>Document</title </head> <body> <script src="https://cdn.jsdelivr.net/npm/@tensorflow/tfjs@2.0.0/dist/tf.min.js"></script> <script type="module" src="./test.js"></script> </body> </html> test.js import * as tf from "./node_modules/@tensorflow/tfjs"; import {loadGraphModel} from "./node_modules/@tensorflow/tfjs-converter"; const MODEL_URL = './model.json'; const model = await loadGraphModel(MODEL_URL); const cat = document.getElementById('cat'); model.execute(tf.browser.fromPixels(cat)); Besides, I run the server using python -m http.server in my command prompt(Windows 10), and this is the error prompt in the console log of my browser: Failed to loa...