mirror of
https://github.com/Watchful1/PushshiftDumps.git
synced 2025-07-25 15:45:19 -04:00
Fix to_csv chunk sizing
This commit is contained in:
parent
33b5b938c1
commit
4e8d6c9b6b
1 changed files with 14 additions and 1 deletions
|
@ -21,12 +21,25 @@ log.setLevel(logging.DEBUG)
|
||||||
log.addHandler(logging.StreamHandler())
|
log.addHandler(logging.StreamHandler())
|
||||||
|
|
||||||
|
|
||||||
|
def read_and_decode(reader, chunk_size, max_window_size, previous_chunk=None, bytes_read=0):
|
||||||
|
chunk = reader.read(chunk_size)
|
||||||
|
bytes_read += chunk_size
|
||||||
|
if previous_chunk is not None:
|
||||||
|
chunk = previous_chunk + chunk
|
||||||
|
try:
|
||||||
|
return chunk.decode()
|
||||||
|
except UnicodeDecodeError:
|
||||||
|
if bytes_read > max_window_size:
|
||||||
|
raise UnicodeError(f"Unable to decode frame after reading {bytes_read:,} bytes")
|
||||||
|
return read_and_decode(reader, chunk_size, max_window_size, chunk, bytes_read)
|
||||||
|
|
||||||
|
|
||||||
def read_lines_zst(file_name):
|
def read_lines_zst(file_name):
|
||||||
with open(file_name, 'rb') as file_handle:
|
with open(file_name, 'rb') as file_handle:
|
||||||
buffer = ''
|
buffer = ''
|
||||||
reader = zstandard.ZstdDecompressor(max_window_size=2**31).stream_reader(file_handle)
|
reader = zstandard.ZstdDecompressor(max_window_size=2**31).stream_reader(file_handle)
|
||||||
while True:
|
while True:
|
||||||
chunk = reader.read(2**27).decode()
|
chunk = read_and_decode(reader, 2**27, (2**29) * 2)
|
||||||
if not chunk:
|
if not chunk:
|
||||||
break
|
break
|
||||||
lines = (buffer + chunk).split("\n")
|
lines = (buffer + chunk).split("\n")
|
||||||
|
|
Loading…
Add table
Add a link
Reference in a new issue