Hello guys, I am building a connection from the mi...
# replication-troubleshooting
g
Hello guys, I am building a connection from the minimal Python Airbyte Source Connector (created with the generator), but the AirbyteRecordMessage requires a valid JSON data as argument. However the data I am collecting is quite large and it is getting complicated to put it in a JSON records format. Is there another way to do this? Maybe collect the whole file and not read its contents?
m
Are you breaking the data into records, usually API returns records inside the complete data. Example:
Copy code
"data": [{"user_id": 1, "name": "marcos"}, {"user_id": 2, "name": "gustavo"}]
g
Yes I am. The problem is not with JSON formatting, it is with the size of the data that does not fit into memory. The process is breaking when I try to transform the pandas dataframe into a JSON
m
You should not read all the dataframe before you yield (just confirming if you are not doing that). Also you should try to do splices so that you read your data in chunks
g
What is you suggestion to read the dataframe in chunks?
m
I thought something similar as the default python http cdk, but I don't think I was on the right path.
The way airbyte works is that the output records are sent to stdout, and the destination read them from stdin. So the best way to read anything is to continually output data, and not hold it in memory. On python we usually do that with a for that yield records, so that the iterable can continually be fed, and after that the record is collected by the garbage collector.
So, if on your read method you are continually outputing AirbyteRecords, memory should not be a problem, given that the data frame do not grow
g
The problem with this specific case is that I have to download a 7zip file from a ftp server, unzip it, load a txt file as a ";" separated file and then I have the records I need. The best way I figured would work is to read the txt file in chunks, but I am not sure if that is the best option
e
Hi. There are multiple ways to split the data. in the api request use stream=True and then iterate the data for example:
Copy code
import requests

s = requests.Session()

list_of_dfs=[]
with s.get(url, headers=None, stream=True) as resp:
  for line in resp.iter_lines():
     if line:
        list_of_dfs.append(pd.DataFrame(line))
Just make sure to map the data acordingly since streaming data is not row by row as you would expect. or use iter_content (iter content also receives param chunk size) and with data you have a buffer that you can put into pd.read_csv and then use chunks to split even more the data you can also use concurrent futures multithreading to create faster stream. This way you can then use pd.concat([list_of_dfs]) after all data is ready ( if needed ofcourse) or if this is big data you can use dask read csv and create partitions or pyspark for the same use cases. if the separator or delimiters are not consistent you will need to work around the data for each stream and create try except block to create a dataframe with the correct data. hope this helps
u
Hello Gustavo Maia, it's been a while without an update from us. Are you still having problems or did you find a solution?
g
I did find a solution. Thanks guys! Reading the file in chunks using pandas solved the memory problems.
s
Happy to hear you found a solution! Let us know if you have more questions by posting in the channel.