1#!/usr/bin/env cwl-runner
2### Pipeline to aggregate data from Climatology Lab
3# Copyright (c) 2021-2022. Harvard University
4#
5# Developed by Research Software Engineering,
6# Faculty of Arts and Sciences, Research Computing (FAS RC)
7# Author: Michael A Bouzinier
8#
9# Licensed under the Apache License, Version 2.0 (the "License");
10# you may not use this file except in compliance with the License.
11# You may obtain a copy of the License at
12#
13# http://www.apache.org/licenses/LICENSE-2.0
14#
15# Unless required by applicable law or agreed to in writing, software
16# distributed under the License is distributed on an "AS IS" BASIS,
17# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
18# See the License for the specific language governing permissions and
19# limitations under the License.
20#
21
22cwlVersion: v1.2
23class: Workflow
24
25requirements:
26 SubworkflowFeatureRequirement: {}
27 StepInputExpressionRequirement: {}
28 InlineJavascriptRequirement: {}
29 ScatterFeatureRequirement: {}
30 MultipleInputFeatureRequirement: {}
31 NetworkAccess:
32 networkAccess: True
33
34
35doc: |
36 This workflow downloads NetCDF datasets from
37 [University of Idaho Gridded Surface Meteorological Dataset](https://www.northwestknowledge.net/metdata/data/),
38 aggregates gridded data to daily mean values over chosen geographies
39 and optionally ingests it into the database.
40
41 The output of the workflow are gzipped CSV files containing
42 aggregated data.
43
44 Optionally, the aggregated data can be ingested into a database
45 specified in the connection parameters:
46
47 * `database.ini` file containing connection descriptions
48 * `connection_name` a string referring to a section in the `database.ini`
49 file, identifying specific connection to be used.
50
51 The workflow can be invoked either by providing command line options
52 as in the following example:
53
54 toil-cwl-runner --retryCount 1 --cleanWorkDir never \
55 --outdir /scratch/work/exposures/outputs \
56 --workDir /scratch/work/exposures \
57 gridmet.cwl \
58 --database /opt/local/database.ini \
59 --connection_name dorieh \
60 --bands rmin rmax \
61 --strategy auto \
62 --geography zcta \
63 --ram 8GB
64
65 Or, by providing a YaML file (see [example](jobs/test_gridmet_job))
66 with similar options:
67
68 toil-cwl-runner --retryCount 1 --cleanWorkDir never \
69 --outdir /scratch/work/exposures/outputs \
70 --workDir /scratch/work/exposures \
71 gridmet.cwl test_gridmet_job.yml
72
73
74inputs:
75 proxy:
76 type: string?
77 default: ""
78 doc: HTTP/HTTPS Proxy if required
79 shapes:
80 type: Directory?
81 doc: Do we even need this parameter, as we instead downloading shapes?
82 geography:
83 type: string
84 doc: |
85 Type of geography: zip codes or counties
86 Valid values: "zip", "zcta" or "county"
87 years:
88 type: string[]
89 default: ['1999', '2000', '2001', '2002', '2003', '2004', '2005', '2006', '2007', '2008', '2009', '2010', '2011', '2012', '2013', '2014', '2015', '2016', '2017', '2018', '2019', '2020']
90 bands:
91 doc: |
92 University of Idaho Gridded Surface Meteorological Dataset
93 [bands](https://developers.google.com/earth-engine/datasets/catalog/IDAHO_EPSCOR_GRIDMET#bands)
94 type: string[]
95 # default: ['bi', 'erc', 'etr', 'fm100', 'fm1000', 'pet', 'pr', 'rmax', 'rmin', 'sph', 'srad', 'th', 'tmmn', 'tmmx', 'vpd', 'vs']
96 strategy:
97 type: string
98 default: auto
99 doc: |
100 [Rasterization strategy](https://foromeplatform.github.io/dorieh/strategy.html)
101 used for spatial aggregation
102 ram:
103 type: string
104 default: 2GB
105 doc: |
106 Runtime memory, available to the process. When aggregation
107 strategy is `auto`, this value is used to calculate the optimal
108 downscaling factor for the available resources.
109 database:
110 type: File
111 doc: Path to database connection file, usually database.ini
112 connection_name:
113 type: string
114 doc: The name of the section in the database.ini file
115 dates:
116 type: string?
117 doc: 'dates restriction, for testing purposes only'
118 domain:
119 type: string
120 default: climate
121
122
123steps:
124 initdb:
125 run: initdb.cwl
126 doc: Ensure that database utilities are at their latest version
127 in:
128 database: database
129 connection_name: connection_name
130 out:
131 - log
132 - err
133
134 init_db_schema:
135 doc: We need to do it because of parallel creation of tables
136 run:
137 class: CommandLineTool
138 baseCommand: [python, -m, dorieh.platform.util.psql]
139 doc: |
140 This tool executes an SQL statement in the database to grant
141 read privileges to NSAPH users (memebrs of group nsaph_admin)
142 inputs:
143 database:
144 type: File
145 doc: Path to database connection file, usually database.ini
146 inputBinding:
147 prefix: --db
148 connection_name:
149 type: string
150 doc: The name of the section in the database.ini file
151 inputBinding:
152 prefix: --connection
153 domain:
154 type: string
155 #default: climate
156 arguments:
157 - valueFrom: $("CREATE SCHEMA IF NOT EXISTS " + inputs.domain + ';')
158 position: 3
159 outputs:
160 log:
161 type: stdout
162 err:
163 type: stderr
164 stderr: "schema.err"
165 stdout: "schema.log"
166 in:
167 database: database
168 connection_name: connection_name
169 domain: domain
170 out:
171 - log
172 - err
173
174 make_registry:
175 run: build_gridmet_model.cwl
176 doc: Writes down YAML file with the database model
177 in:
178 depends_on: init_db_schema/log
179 domain: domain
180 out:
181 - model
182 - log
183 - errors
184
185 init_tables:
186 doc: creates or recreates database tables, one for each band
187 scatter:
188 - band
189 run:
190 class: Workflow
191 inputs:
192 depends_on:
193 type: Any?
194 doc: a special field used to enforce dependencies and execution order
195 registry:
196 type: File
197 table:
198 type: string
199 domain:
200 type: string
201 database:
202 type: File
203 connection_name:
204 type: string
205 steps:
206 reset:
207 run: reset.cwl
208 in:
209 registry: registry
210 domain: domain
211 database: database
212 connection_name: connection_name
213 table: table
214 out:
215 - log
216 - errors
217 index:
218 run: index.cwl
219 in:
220 depends_on: reset/log
221 registry: registry
222 domain: domain
223 table: table
224 database: database
225 connection_name: connection_name
226 out: [log, errors]
227 outputs:
228 reset_log:
229 type: File
230 outputSource: reset/log
231 reset_err:
232 type: File
233 outputSource: reset/errors
234 index_log:
235 type: File
236 outputSource: index/log
237 index_err:
238 type: File
239 outputSource: index/errors
240 in:
241 depends_on: initdb/log
242 registry: make_registry/model
243 database: database
244 connection_name: connection_name
245 band: bands
246 geography: geography
247 domain: domain
248 table:
249 valueFrom: $(inputs.geography + '_' + inputs.band)
250 out:
251 - reset_log
252 - reset_err
253 - index_log
254 - index_err
255
256 process:
257 run: gridmet_one_file.cwl
258 doc: Downloads raw data and aggregates it over shapes and time
259 scatter:
260 - band
261 - year
262 scatterMethod: nested_crossproduct
263
264 in:
265 proxy: proxy
266 depends_on: init_tables/index_log
267 model: make_registry/model
268 shapes: shapes
269 geography: geography
270 strategy: strategy
271 ram: ram
272 year: years
273 dates: dates
274 band: bands
275 database: database
276 connection_name: connection_name
277 domain: domain
278 months:
279 valueFrom: $([1,2,3,4,5,6,7,8,9,10,11,12])
280 table:
281 valueFrom: $(inputs.geography + '_' + inputs.band)
282
283 out:
284 - download_log
285 - download_err
286 - add_data_aggregate_errors
287 - add_data_data
288 - add_data_aggregate_log
289 - add_data_ingest_log
290 - add_data_ingest_errors
291 - vacuum_log
292 - vacuum_err
293
294 export:
295 run: export.cwl
296 scatter:
297 - band
298 in:
299 depends_on: process/vacuum_log
300 database: database
301 connection_name: connection_name
302 format:
303 valueFrom: "parquet"
304 domain: domain
305 geography: geography
306 band: bands
307 table:
308 valueFrom: $(inputs.domain + '.' + inputs.geography + '_' + inputs.band)
309 partition:
310 valueFrom: $(["year"])
311 output:
312 valueFrom: $('export/' + inputs.domain + '/' + inputs.geography + '_' + inputs.band)
313 out:
314 - data
315 - log
316 - errors
317
318
319
320outputs:
321 initdb_log:
322 type: File?
323 outputSource: initdb/log
324 initdb_err:
325 type: File?
326 outputSource: initdb/err
327
328 init_schema_log:
329 type: File?
330 outputSource: init_db_schema/log
331 init_schema_err:
332 type: File?
333 outputSource: init_db_schema/err
334
335 registry:
336 type: File?
337 outputSource: make_registry/model
338 registry_log:
339 type: File?
340 outputSource: make_registry/log
341 registry_err:
342 type: File?
343 outputSource: make_registry/errors
344
345 data:
346 type:
347 type: array
348 items:
349 type: array
350 items:
351 type: array
352 items: [File]
353 outputSource: process/add_data_data
354 download_log:
355 type:
356 type: array
357 items:
358 type: array
359 items: [File]
360 outputSource: process/download_log
361 download_err:
362 type:
363 type: array
364 items:
365 type: array
366 items: [File]
367 outputSource: process/download_err
368
369 process_log:
370 type:
371 type: array
372 items:
373 type: array
374 items:
375 type: array
376 items: [File]
377 outputSource: process/add_data_aggregate_log
378 process_err:
379 type:
380 type: array
381 items:
382 type: array
383 items:
384 type: array
385 items: [File]
386 outputSource: process/add_data_aggregate_errors
387
388 ingest_log:
389 type:
390 type: array
391 items:
392 type: array
393 items:
394 type: array
395 items: [File]
396 outputSource: process/add_data_ingest_log
397 ingest_err:
398 type:
399 type: array
400 items:
401 type: array
402 items:
403 type: array
404 items: [File]
405 outputSource: process/add_data_ingest_errors
406
407 reset_log:
408 type:
409 type: array
410 items: [File]
411 outputSource: init_tables/reset_log
412 reset_err:
413 type:
414 type: array
415 items: [File]
416 outputSource: init_tables/reset_err
417
418 index_log:
419 type:
420 type: array
421 items: [File]
422 outputSource: init_tables/index_log
423 index_err:
424 type:
425 type: array
426 items: [File]
427 outputSource: init_tables/index_err
428
429 vacuum_log:
430 type:
431 type: array
432 items:
433 type: array
434 items: [File]
435 outputSource: process/vacuum_log
436 vacuum_err:
437 type:
438 type: array
439 items:
440 type: array
441 items: [File]
442 outputSource: process/vacuum_err
443
444 export_data:
445 type:
446 type: array
447 items: ['File', 'Directory']
448 outputSource: export/data
449 export_log:
450 type: File[]
451 outputSource: export/log
452 export_err:
453 type: File[]
454 outputSource: export/errors