gridmet.cwl

  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