noaa-goes-visualization/puller_fits.py

258 lines
12 KiB
Python
Raw Normal View History

2024-05-16 15:43:12 -04:00
import os
import urllib.request
import urllib.parse
import re
import time
import random
import datetime
import json
from threading import Thread
import queue
import math
from functools import partial
2024-07-05 10:24:41 -04:00
from queue import Empty
2024-05-16 15:43:12 -04:00
import tqdm
directory_url = r"https://data.ngdc.noaa.gov/platforms/solar-space-observing-satellites/goes/"
stored_images_dir = r"..\Data"
# directory_url = r"https://data.ngdc.noaa.gov/platforms/solar-space-observing-satellites/goes/goes16/l2/data/ephe-l2-orb1m/2017/12/"
# stored_images_dir = r"Z:\NOAA GOES Data\Data\goes16/l2/data/ephe-l2-orb1m/2017/12/"
ignore_folder_names = ["l1b", "goes17", "2017", "2018", "2019", "2020", "2021", "2022", "2023"]
2024-05-16 15:43:12 -04:00
file_database_path = r"..\file_database.json"
fetch_interval = 0 # 60*60
nfetchworkers = 10 # Be nice to the servers, this value is how many threads will be asking for links and file info at the same time
ndownloadworkers = 4 # Be nice to the servers, this value is how many threads will be downloading files at the same time
randomize_order = False
2024-05-16 15:43:12 -04:00
links_regex_pattern = r'(?<=<a href=")([^ ?:]*)(?=">.*\d{4}-\d{2}-\d{2} \d{2}:\d{2})' # Find href links that do not contain question marks or whitespace
times_regex_pattern = r'(?<=<td align="right">)(\d{4}-\d{2}-\d{2} \d{2}:\d{2})(?= )' # Find timestamps in UTC in the format YYYY-MM-DD HH:mm
sizes_regex_pattern = r'(?<=\d{4}-\d{2}-\d{2} \d{2}:\d{2} <\/td><td align="right">)([ .\d]{1,3})([KMG]|- |0 | )(?=<\/td>)' # Find file size markers in one of many possible formats
links_matcher = re.compile(links_regex_pattern)
times_matcher = re.compile(times_regex_pattern)
sizes_matcher = re.compile(sizes_regex_pattern)
2024-05-16 15:43:12 -04:00
def find_links_worker(query_work_queue, query_result_queue, attempt_count = 3, randomize_order = True):
while True:
job = query_work_queue.get()
if job == None:
return
url= job
attempts = 0
html_content = None
while attempts < attempt_count:
try:
with urllib.request.urlopen(url) as response:
html_content = response.read().decode('utf-8')
break
except Exception as e:
print(f"Exception while fetching links: {e}")
time.sleep(1 + random.random())
attempts += 1
if (attempt_count > 1) and (attempts == attempt_count):
print(f'\nAfter {attempt_count} retries, could not fetch: {url}')
return
links = links_matcher.findall(html_content)
times = times_matcher.findall(html_content)
sizes = sizes_matcher.findall(html_content)
for i in range(len(times)):
dt = datetime.datetime.strptime(times[i].strip(), "%Y-%m-%d %H:%M")
dt.replace(tzinfo=datetime.timezone.utc)
times[i] = time.mktime(dt.timetuple())
for i in range(len(sizes)):
match sizes[i][1].strip():
case '-':
sizes[i] = 0
case 'K':
sizes[i] = int(float(sizes[i][0].strip()) * 1024)
case 'M':
sizes[i] = int(float(sizes[i][0].strip()) * 1024 * 1024)
case 'G':
sizes[i] = int(float(sizes[i][0].strip()) * 1024 * 1024 * 1024)
case "0":
sizes[i] = 0
case "":
sizes[i] = int(float(sizes[i][0].strip()))
case _:
raise(ValueError(f"Unexpected symbol while parsing links page: {_}"))
if (len(links) != len(times)) or (len(times) != len(sizes)):
raise(ValueError("Links parsing error, mismatched numbers of links, times, or sizes!"))
results = list(zip(links, times, sizes))
if randomize_order:
random.shuffle(results)
for _link, _time, _size in results:
if _link.endswith("/"):
if _link.split(r"/")[-2] in ignore_folder_names:
continue
else:
query_work_queue.put(urllib.parse.urljoin(url,_link))
2024-05-16 15:43:12 -04:00
else:
query_result_queue.put((urllib.parse.urljoin(url,_link), _time, _size))
2024-05-16 15:43:12 -04:00
# Fetch file from url and store it to path, retrying on failure
def file_download_worker(download_work_queue, download_result_queue, attempt_count = 2):
2024-05-16 15:43:12 -04:00
while True:
job = download_work_queue.get()
2024-05-16 15:43:12 -04:00
if job == None:
return
url, path, t, s = job
2024-05-16 15:43:12 -04:00
attempts = 0
while attempts < attempt_count:
try:
req = urllib.request.Request(url, data=None)
image_data = urllib.request.urlopen(req).read()
if math.isclose(len(image_data), s, rel_tol=0.05):
os.makedirs(os.path.split(path)[0], exist_ok=True)
open(path, 'wb').write(image_data)
download_result_queue.put((True, url, t))
break
else:
raise ValueError("Downloaded file is the wrong size!")
2024-05-16 15:43:12 -04:00
except Exception as e:
if hasattr(e, "code") and e.code == 404: # This is expected if the file has been removed from the site (at least for swpc.noaa.gov)
attempts += attempt_count
elif attempts == 0:
print(f"\nA problem occurred on image: {url} | {e}")
time.sleep(1 + random.random())
attempts += 1
if (attempt_count > 1) and (attempts == attempt_count):
print(f'\nAfter {attempt_count} retries, could not fetch: {url}')
download_result_queue.put((False, url, t))
2024-05-16 15:43:12 -04:00
break
if __name__ == "__main__":
file_info_cache = {}
try:
print(f"Attempting to load file records from cache: {file_database_path}")
with open(file_database_path, 'r') as f:
file_info_cache = json.loads(f.read())
print(f"File records loaded from cache: {len(file_info_cache)} records found.")
except Exception as e:
print(f"Load failed, starting with empty cache")
file_info_cache = {}
query_work_queue = queue.Queue()
query_result_queue = queue.Queue()
download_work_queue = queue.Queue()
download_result_queue = queue.Queue()
2024-05-16 15:43:12 -04:00
workers = []
for _ in range(nfetchworkers):
t = Thread(target=find_links_worker, args=(query_work_queue, query_result_queue, 3, randomize_order), daemon=True)
t.start()
workers.append(t)
for _ in range(ndownloadworkers):
t = Thread(target=file_download_worker, args=(download_work_queue, download_result_queue), daemon=True)
2024-05-16 15:43:12 -04:00
t.start()
workers.append(t)
2024-05-16 15:43:12 -04:00
try:
while True:
fetched_image_count = 0
already_had_image_count = 0
failed_image_count = 0
urllen = len(directory_url)
query_work_queue.put(directory_url)
for l, t, s in tqdm.tqdm(iter(partial(query_result_queue.get, timeout=30.0), None), desc="Downloading files"):
2024-05-16 15:43:12 -04:00
# Collect completed jobs and record completion status
while True:
try:
r_success, r_url, r_t = download_result_queue.get_nowait()
2024-05-16 15:43:12 -04:00
if r_success:
fetched_image_count += 1
file_info_cache[r_url] = r_t
else:
failed_image_count += 1
except queue.Empty:
break
file_portion_of_link = l[urllen:]
filepath = os.path.join(stored_images_dir, file_portion_of_link.replace("/", os.sep))
2024-05-16 15:43:12 -04:00
# If we dont have the file or the file at the link is newer than the one we previously fetched
if (not (l in file_info_cache)) or t > file_info_cache[l]:
if filepath.endswith(".fits"):
filepath2 = filepath.split(".fits")[0] + "_f.fits" # also look for the filtered version of the file
filepath3 = filepath.split(".fits")[0] + "_e.fits" # also look for the error version of the file
if os.path.exists(filepath2) or os.path.exists(filepath3): # If the unfiltered filename exists, that case will be handled in the alternative code path in if os.path.exists(filepath):
if l in file_info_cache: # If we have record of this file, it must be out of date, delete it and download the new version.
try:
os.remove(filepath2)
os.remove(filepath3)
except:
pass
else: # If we have no record of this file, update the file info cache and don't redownload
file_info_cache[l] = t
already_had_image_count += 1
continue
2024-05-16 15:43:12 -04:00
if os.path.exists(filepath):
if l in file_info_cache: # If we have record of this file, it must be out of date, delete it and download the new version.
if os.path.exists(filepath):
os.remove(filepath)
else: # If we have no record of this file, update the file info cache and don't redownload if fsize is right
fsize = os.path.getsize(filepath)
if math.isclose(fsize, s, rel_tol=0.05):
file_info_cache[l] = t
already_had_image_count += 1
continue
else:
print(f'Found a mismatched size on file: {filepath} Redownloading!')
download_work_queue.put((l, filepath, t, s))
2024-05-16 15:43:12 -04:00
else:
# We have a download record, confirm the file actually exists on disk
if os.path.exists(filepath):
already_had_image_count += 1
elif filepath.endswith(".fits"):
filepath2 = filepath.split(".fits")[0] + "_f.fits" # also look for the filtered version of the file
filepath3 = filepath.split(".fits")[0] + "_e.fits" # also look for the error version of the file
if os.path.exists(filepath2) or os.path.exists(filepath3):
already_had_image_count += 1
else: # We could not find the file on disk, queue for redownload
download_work_queue.put((l, filepath, t))
2024-05-16 15:43:12 -04:00
with open(file_database_path, 'w') as f:
f.write(json.dumps(file_info_cache))
print(f"Downloaded {fetched_image_count} | Already had {already_had_image_count} | Failed {failed_image_count}")
if fetch_interval > 0:
print(f"Run complete!, sleeping for {fetch_interval} seconds.")
with open(file_database_path, 'w') as f:
f.write(json.dumps(file_info_cache))
time.sleep(fetch_interval)
else:
print("Run complete!, exiting...")
break
2024-05-16 15:43:12 -04:00
except KeyboardInterrupt:
print("Saving file database and shutting down.")
2024-07-05 10:24:41 -04:00
except Empty:
print("Work Complete, Shutting down...")
except Exception as e:
print(f"Unhandled Exception during run: {e}")
print("Shutting down")
2024-07-05 10:24:41 -04:00
with open(file_database_path, 'w') as f:
f.write(json.dumps(file_info_cache))
for _ in range(nfetchworkers):
try:
query_work_queue.put(None, timeout=5.0)
except:
break
2024-05-16 15:43:12 -04:00
2024-07-05 10:24:41 -04:00
time.sleep(1)
print(f"Downloaded {fetched_image_count} | Already had {already_had_image_count} | Failed {failed_image_count}")
for _ in range(ndownloadworkers):
2024-05-16 15:43:12 -04:00
try:
download_work_queue.put(None, timeout=5.0)
2024-05-16 15:43:12 -04:00
except:
break
for w in workers:
2024-07-05 10:24:41 -04:00
w.join(5.0)