-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathharvester.rb
More file actions
159 lines (129 loc) · 4.39 KB
/
Copy pathharvester.rb
File metadata and controls
159 lines (129 loc) · 4.39 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
require 'net/http'
require 'json'
require 'time'
require 'fileutils'
require 'securerandom'
require_relative 'mapper'
class Harvester
HDX_API_URL = "https://data.humdata.org/api/3/action/package_search"
BATCH_SIZE = 1000
# Default paths for production datastores
STATE_FILE = "state.json"
HDX_METADATA_DIR = "metadata-hdx"
AARDVARK_METADATA_DIR = "metadata-aardvark"
def self.run(state_file: STATE_FILE, hdx_metadata_dir: HDX_METADATA_DIR, aardvark_metadata_dir: AARDVARK_METADATA_DIR)
new(
state_file: state_file,
hdx_metadata_dir: hdx_metadata_dir,
aardvark_metadata_dir: aardvark_metadata_dir
).run
end
attr_reader :state_file, :hdx_metadata_dir, :aardvark_metadata_dir
def initialize(state_file: STATE_FILE, hdx_metadata_dir: HDX_METADATA_DIR, aardvark_metadata_dir: AARDVARK_METADATA_DIR)
@state_file = state_file
@hdx_metadata_dir = hdx_metadata_dir
@aardvark_metadata_dir = aardvark_metadata_dir
end
def run
last_run = load_last_run
puts "Last run date: #{last_run}"
# Prepare the fq filter
last_run_str = last_run.strftime('%Y-%m-%dT%H:%M:%SZ')
fq = "metadata_modified:[#{last_run_str} TO *]"
start = 0
total_count = 0
total_fetched_count = 0
processed_count = 0
current_time = Time.now.utc
loop do
batch, count = fetch_datasets(fq: fq, start: start)
break if batch.empty?
total_count = count
page_number = (start / BATCH_SIZE) + 1
puts "Fetching page #{page_number} of #{total_count}"
total_fetched_count += batch.size
processed_count += process_batch(batch, last_run)
start += BATCH_SIZE
break if start >= total_count
end
puts "Fetched #{total_fetched_count} datasets from HDX."
exit if total_fetched_count == 0
puts "Processed #{processed_count} new datasets."
update_last_run(current_time)
end
private
# Process a batch of datasets and return the count of successfully processed new datasets
def process_batch(batch, last_run)
processed = 0
batch.each do |dataset|
processed += 1 if process_dataset(dataset, last_run)
end
processed
end
# Process, map, and save a single dataset if it was modified since the last run
def process_dataset(dataset, last_run)
# Filter by modified date
modified_date_str = dataset['metadata_modified']
return false unless modified_date_str
modified_date = Time.parse(modified_date_str).utc
return false unless modified_date > last_run
# Use the ID for file naming
id = dataset['id'] || dataset['name'] || "unknown_#{SecureRandom.hex(4)}"
# Save original
save_metadata(@hdx_metadata_dir, id, dataset)
# Map and save Aardvark
aardvark_data = Mapper.map(dataset)
save_metadata(@aardvark_metadata_dir, id, aardvark_data)
puts "Processed #{id}"
true
end
def load_last_run
if File.exist?(@state_file)
data = JSON.parse(File.read(@state_file))
Time.parse(data['last_run']).utc
else
Time.parse("2024-01-01T00:00:00Z").utc
end
rescue StandardError => e
puts "Error loading state file: #{e.message}"
Time.parse("2024-01-01T00:00:00Z").utc
end
def update_last_run(time)
File.write(@state_file, JSON.pretty_generate({ last_run: time.strftime('%Y-%m-%dT%H:%M:%SZ') }))
puts "Updated state file with: #{time.strftime('%Y-%m-%dT%H:%M:%SZ')}"
rescue StandardError => e
puts "Error updating state file: #{e.message}"
end
def fetch_datasets(fq: nil, start: 0)
params = {
"q" => "has_geodata:true",
"rows" => BATCH_SIZE,
}
params["start"] = start
params["fq"] = fq if fq
uri = URI(HDX_API_URL)
uri.query = URI.encode_www_form(params)
puts "fetching from #{uri}"
response = Net::HTTP.get_response(uri)
if response.is_a?(Net::HTTPSuccess)
data = JSON.parse(response.body)
return data.dig('result', 'results') || [], data.dig('result', 'count') || 0
else
puts "Error fetching from HDX API: #{response.code} #{response.message}"
return [], 0
end
rescue StandardError => e
puts "Request error: #{e.message}"
return [], 0
end
def save_metadata(dir, id, data)
FileUtils.mkdir_p(dir)
filename = File.join(dir, "#{id}.json")
File.write(filename, JSON.pretty_generate(data))
rescue StandardError => e
puts "Error saving metadata for #{id}: #{e.message}"
end
end
if __FILE__ == $0
Harvester.run
end