openrvdas.logger.transforms.true_winds_transform
Compute true winds by processing and aggregating vessel course/speed/heading and relative wind dir/speed records.
There are plenty of challenges with computing a universally-accepted true wind value. Even with the correct algorithm (not a given), unless the vessel nav and anemometer values have identical timestamps, there's the question of how one integrates/interpolates/extrapolates values with different timestamps.
The update_on_fields dict allows specifying which fields will trigger a new computation. This allows making simplifying assumptions, e.g. that vessel course/speed/heading is less variable than wind dir/speed, so we will only produce updated results when we receive new anemometer records. The max_field_age dict allows specifying a time beyond which we are unwilling to trust any computations.
A more robust approach would be to wait until we got the next vessel record and interpolate the course/speed/heading values between the two vessel records (or, conversely, output when we got a vessel record, using an interpolation of the preceding and following anemometer readings?).
1#!/usr/bin/env python3 2"""Compute true winds by processing and aggregating vessel 3course/speed/heading and relative wind dir/speed records. 4 5There are plenty of challenges with computing a universally-accepted 6true wind value. Even with the correct algorithm (not a given), unless 7the vessel nav and anemometer values have identical timestamps, 8there's the question of how one integrates/interpolates/extrapolates 9values with different timestamps. 10 11The update_on_fields dict allows specifying which fields will trigger 12a new computation. This allows making simplifying assumptions, e.g. that 13vessel course/speed/heading is less variable than wind dir/speed, so we 14will only produce updated results when we receive new anemometer records. 15The max_field_age dict allows specifying a time beyond which we are unwilling 16to trust any computations. 17 18A more robust approach would be to wait until we got the next vessel 19record and interpolate the course/speed/heading values between the two 20vessel records (or, conversely, output when we got a vessel record, 21using an interpolation of the preceding and following anemometer 22readings?). 23 24""" 25 26import logging 27import time 28 29from pprint import pformat 30from typing import Union 31 32from logger.utils.das_record import DASRecord, to_das_record_list # noqa: E402 33from logger.utils.truewinds.truew import truew # noqa: E402 34from logger.transforms.derived_data_transform import DerivedDataTransform # noqa: E402 35 36 37################################################################################ 38class TrueWindsTransform(DerivedDataTransform): 39 """Transform that computes true winds from vessel 40 course/speed/heading and anemometer relative wind speed/dir. 41 """ 42 43 def __init__(self, 44 course_field, speed_field, heading_field, 45 wind_dir_field, wind_speed_field, 46 true_dir_name, 47 true_speed_name, 48 apparent_dir_name, 49 update_on_fields=[], 50 max_field_age={}, 51 zero_line_reference=0, 52 convert_wind_factor=1, 53 convert_speed_factor=1, 54 data_id=None, 55 metadata_interval=None, 56 **kwargs): 57 """ 58 ``` 59 course_field 60 speed_field 61 heading_field 62 wind_dir_field 63 wind_speed_field 64 Field names from which we should take values for 65 course, speed over ground, heading, relative wind speed 66 and relative wind direction. 67 68 true_dir_name 69 true_speed_name 70 apparent_dir_name 71 Names that should be given to transform output values. 72 73 update_on_fields 74 If non-empty, a list of fields, any of whose arrival should 75 trigger an output record. If None, generate output when any 76 field is updated. 77 78 max_field_age 79 If non-empty, a dict of field_name:seconds, specifying that 80 no output is to be produces if the age of any of the specified 81 names is older than the specified number of seconds. 82 83 zero_line_reference 84 Angle between bow and zero line on anemometer, referenced 85 to ship. 86 87 convert_wind_factor 88 convert_speed_factor 89 Wind speed and vessel speed may be in different units; e.g 90 wind speed in meters/sec and vessel speed in knots. Multiply 91 the respective values we get by the respective factors. 92 Typically, only one of these will be not equal to 1; e.g. to 93 output true winds as meters/sec, we'll leave convert_wind_factor 94 as 1 and specify convert_speed_factor=0.5144 95 96 metadata_interval - how many seconds between when we attach field metadata 97 to a record we send out. 98 99 data_id Optional name that will be attached to the resulting DASRecord 100 ``` 101 """ 102 super().__init__(**kwargs) # processes 'quiet' and type hints 103 104 self.course_field = course_field 105 self.speed_field = speed_field 106 self.heading_field = heading_field 107 self.wind_dir_field = wind_dir_field 108 self.wind_speed_field = wind_speed_field 109 110 self.true_dir_name = true_dir_name 111 self.true_speed_name = true_speed_name 112 self.apparent_dir_name = apparent_dir_name 113 114 self.update_on_fields = update_on_fields 115 self.max_field_age = max_field_age 116 self.field_age = {} 117 118 self.zero_line_reference = zero_line_reference 119 120 self.convert_wind_factor = convert_wind_factor 121 self.convert_speed_factor = convert_speed_factor 122 123 self.metadata_interval = metadata_interval 124 self.last_metadata_send = 0 125 126 self.data_id = data_id 127 128 # TODO: It may make sense for us to cache most recent values so 129 # that, for example, we can take single DASRecords in the 130 # transform() method and use the most recent values we've seen 131 # from previous calls. 132 self.course_val = None 133 self.speed_val = None 134 self.heading_val = None 135 self.wind_dir_val = None 136 self.wind_speed_val = None 137 138 self.course_val_time = 0 139 self.speed_val_time = 0 140 self.heading_val_time = 0 141 self.wind_dir_val_time = 0 142 self.wind_speed_val_time = 0 143 144 ############################ 145 def fields(self): 146 """Which fields are we interested in to produce transformed data?""" 147 return [self.course_field, self.speed_field, self.heading_field, 148 self.wind_dir_field, self.wind_speed_field] 149 150 ############################ 151 def _metadata(self): 152 """Return a dict of metadata for our derived fields.""" 153 154 metadata_fields = { 155 self.true_dir_name: { 156 'description': 'Derived true wind direction from %s, %s, %s, %s, %s' 157 % (self.course_field, self.speed_field, self.heading_field, 158 self.wind_dir_field, self.wind_speed_field), 159 'units': 'degrees', 160 'device': 'TrueWindTransform', 161 'device_type': 'DerivedTrueWindTransform', 162 'device_type_field': self.true_dir_name 163 }, 164 self.true_speed_name: { 165 'description': 'Derived true wind speed from %s, %s, %s, %s, %s' 166 % (self.course_field, self.speed_field, self.heading_field, 167 self.wind_dir_field, self.wind_speed_field), 168 'units': 'depends on conversion used for %s (%g) and %s (%g)' 169 % (self.speed_field, self.convert_speed_factor, 170 self.wind_speed_field, self.convert_wind_factor), 171 'device': 'TrueWindTransform', 172 'device_type': 'DerivedTrueWindTransform', 173 'device_type_field': self.true_speed_name 174 }, 175 self.apparent_dir_name: { 176 'description': 'Derived apparent wind speed from %s, %s, %s, %s, %s' 177 % (self.course_field, self.speed_field, self.heading_field, 178 self.wind_dir_field, self.wind_speed_field), 179 'units': 'degrees', 180 'device': 'TrueWindTransform', 181 'device_type': 'DerivedTrueWindTransform', 182 'device_type_field': self.apparent_dir_name 183 } 184 } 185 return metadata_fields 186 187 ############################ 188 def transform(self, record: Union[dict, DASRecord]): 189 """Incorporate any useable fields in this record, and if it gives 190 us a new true wind value, return the results.""" 191 192 # See if it's something we can process, and if not, try digesting 193 if not self.can_process_record(record): # inherited from BaseModule() 194 return self.digest_record(record) # inherited from BaseModule() 195 196 results = [] 197 for das_record in to_das_record_list(record): 198 # If they haven't specified specific fields we should wait for 199 # before updates, plan to emit an update after every new record 200 # we process. Otherwise, assume we're not going to update unless 201 # we see one of the named fields. 202 if not self.update_on_fields: 203 update = True 204 else: 205 update = False 206 207 timestamp = das_record.timestamp 208 if not timestamp: 209 logging.info('DASRecord is missing timestamp - skipping') 210 continue 211 212 # Get latest values for any of our fields 213 fields = das_record.fields 214 if self.course_field in fields: 215 if timestamp >= self.course_val_time: 216 self.course_val = fields.get(self.course_field) 217 self.course_val_time = timestamp 218 if self.course_field in self.update_on_fields: 219 update = True 220 221 if self.speed_field in fields: 222 if timestamp >= self.speed_val_time: 223 self.speed_val = fields.get(self.speed_field) 224 self.speed_val *= self.convert_speed_factor 225 self.speed_val_time = timestamp 226 if self.speed_field in self.update_on_fields: 227 update = True 228 229 if self.heading_field in fields: 230 if timestamp >= self.heading_val_time: 231 self.heading_val = fields.get(self.heading_field) 232 self.heading_val_time = timestamp 233 if self.heading_field in self.update_on_fields: 234 update = True 235 236 if self.wind_dir_field in fields: 237 if timestamp >= self.wind_dir_val_time: 238 self.wind_dir_val = fields.get(self.wind_dir_field) 239 self.wind_dir_val_time = timestamp 240 if self.wind_dir_field in self.update_on_fields: 241 update = True 242 243 if self.wind_speed_field in fields: 244 if timestamp >= self.wind_speed_val_time: 245 self.wind_speed_val = fields.get(self.wind_speed_field) 246 self.wind_speed_val *= self.convert_wind_factor 247 self.wind_speed_val_time = timestamp 248 if self.wind_speed_field in self.update_on_fields: 249 update = True 250 251 # Check if needed all values are present, and none are too old to use 252 if self._values_too_old(timestamp): 253 continue 254 255 # If we've not seen anything that updates fields that would 256 # trigger a new true winds value, skip rest of computation. 257 if not update: 258 logging.debug('No update needed') 259 continue 260 261 logging.debug('Computing new true winds') 262 (true_dir, true_speed, apparent_dir) = truew(crse=self.course_val, 263 cspd=self.speed_val, 264 hd=self.heading_val, 265 wdir=self.wind_dir_val, 266 zlr=self.zero_line_reference, 267 wspd=self.wind_speed_val) 268 269 logging.debug('Got true winds: dir: %s, speed: %s, apparent_dir: %s', 270 true_dir, true_speed, apparent_dir) 271 if None in (true_dir, true_speed, apparent_dir): 272 logging.info('Got invalid true winds') 273 continue 274 275 # If here, we've got a valid new true wind result 276 true_wind_fields = {self.true_dir_name: true_dir, 277 self.true_speed_name: true_speed, 278 self.apparent_dir_name: apparent_dir} 279 280 # Add in metadata if so specified and it's been long enough since 281 # we last sent it. 282 now = time.time() 283 if self.metadata_interval and \ 284 now - self.metadata_interval > self.last_metadata_send: 285 metadata = {'fields': self._metadata()} 286 self.last_metadata_send = now 287 logging.debug('Emitting metadata: %s', pformat(metadata)) 288 else: 289 metadata = None 290 291 results.append(DASRecord(timestamp=timestamp, fields=true_wind_fields, 292 metadata=metadata, data_id=self.data_id)) 293 294 logging.debug('Sending %d true wind results.', len(results)) 295 return results 296 297 ############################ 298 def _values_too_old(self, timestamp): 299 """Return true if any values are missing or too old to use.""" 300 301 if None in (self.course_val, self.speed_val, self.heading_val, 302 self.wind_dir_val, self.wind_speed_val): 303 logging.debug('Not all required values for true winds are present: ' 304 'time: %s: %s: %s, %s: %s, %s: %s, %s: %s, %s: %s', 305 timestamp, 306 self.course_field, self.course_val, 307 self.speed_field, self.speed_val, 308 self.heading_field, self.heading_val, 309 self.wind_dir_field, self.wind_dir_val, 310 self.wind_speed_field, self.wind_speed_val) 311 return True 312 313 course_max_age = self.max_field_age.get(self.course_field) 314 if (course_max_age and timestamp - self.course_val_time > course_max_age): 315 logging.debug('course_field too old - max age %g, age %g', 316 course_max_age, timestamp - self.course_val_time) 317 return True 318 319 speed_max_age = self.max_field_age.get(self.speed_field) 320 if speed_max_age: 321 if timestamp - self.speed_val_time > speed_max_age: 322 logging.debug('speed_field too old - max age %g, age %g', 323 speed_max_age, timestamp - self.speed_val_time) 324 return True 325 326 heading_max_age = self.max_field_age.get(self.heading_field) 327 if heading_max_age: 328 if timestamp - self.heading_val_time > heading_max_age: 329 logging.debug('heading_field too old - max age %g, age %g', 330 heading_max_age, timestamp - self.heading_val_time) 331 return True 332 333 wind_dir_max_age = self.max_field_age.get(self.wind_dir_field) 334 if wind_dir_max_age: 335 if timestamp - self.wind_dir_val_time > wind_dir_max_age: 336 logging.debug('wind_dir_field too old - max age %g, age %g', 337 wind_dir_max_age, timestamp - self.wind_dir_val_time) 338 return True 339 340 wind_speed_max_age = self.max_field_age.get(self.wind_speed_field) 341 if wind_speed_max_age: 342 if timestamp - self.wind_speed_val_time > wind_speed_max_age: 343 logging.debug('wind_speed_field too old - max age %g, age %g', 344 wind_speed_max_age, timestamp - self.wind_speed_val_time) 345 return True 346 347 # Everything is present, and nothing's too old... 348 return False
39class TrueWindsTransform(DerivedDataTransform): 40 """Transform that computes true winds from vessel 41 course/speed/heading and anemometer relative wind speed/dir. 42 """ 43 44 def __init__(self, 45 course_field, speed_field, heading_field, 46 wind_dir_field, wind_speed_field, 47 true_dir_name, 48 true_speed_name, 49 apparent_dir_name, 50 update_on_fields=[], 51 max_field_age={}, 52 zero_line_reference=0, 53 convert_wind_factor=1, 54 convert_speed_factor=1, 55 data_id=None, 56 metadata_interval=None, 57 **kwargs): 58 """ 59 ``` 60 course_field 61 speed_field 62 heading_field 63 wind_dir_field 64 wind_speed_field 65 Field names from which we should take values for 66 course, speed over ground, heading, relative wind speed 67 and relative wind direction. 68 69 true_dir_name 70 true_speed_name 71 apparent_dir_name 72 Names that should be given to transform output values. 73 74 update_on_fields 75 If non-empty, a list of fields, any of whose arrival should 76 trigger an output record. If None, generate output when any 77 field is updated. 78 79 max_field_age 80 If non-empty, a dict of field_name:seconds, specifying that 81 no output is to be produces if the age of any of the specified 82 names is older than the specified number of seconds. 83 84 zero_line_reference 85 Angle between bow and zero line on anemometer, referenced 86 to ship. 87 88 convert_wind_factor 89 convert_speed_factor 90 Wind speed and vessel speed may be in different units; e.g 91 wind speed in meters/sec and vessel speed in knots. Multiply 92 the respective values we get by the respective factors. 93 Typically, only one of these will be not equal to 1; e.g. to 94 output true winds as meters/sec, we'll leave convert_wind_factor 95 as 1 and specify convert_speed_factor=0.5144 96 97 metadata_interval - how many seconds between when we attach field metadata 98 to a record we send out. 99 100 data_id Optional name that will be attached to the resulting DASRecord 101 ``` 102 """ 103 super().__init__(**kwargs) # processes 'quiet' and type hints 104 105 self.course_field = course_field 106 self.speed_field = speed_field 107 self.heading_field = heading_field 108 self.wind_dir_field = wind_dir_field 109 self.wind_speed_field = wind_speed_field 110 111 self.true_dir_name = true_dir_name 112 self.true_speed_name = true_speed_name 113 self.apparent_dir_name = apparent_dir_name 114 115 self.update_on_fields = update_on_fields 116 self.max_field_age = max_field_age 117 self.field_age = {} 118 119 self.zero_line_reference = zero_line_reference 120 121 self.convert_wind_factor = convert_wind_factor 122 self.convert_speed_factor = convert_speed_factor 123 124 self.metadata_interval = metadata_interval 125 self.last_metadata_send = 0 126 127 self.data_id = data_id 128 129 # TODO: It may make sense for us to cache most recent values so 130 # that, for example, we can take single DASRecords in the 131 # transform() method and use the most recent values we've seen 132 # from previous calls. 133 self.course_val = None 134 self.speed_val = None 135 self.heading_val = None 136 self.wind_dir_val = None 137 self.wind_speed_val = None 138 139 self.course_val_time = 0 140 self.speed_val_time = 0 141 self.heading_val_time = 0 142 self.wind_dir_val_time = 0 143 self.wind_speed_val_time = 0 144 145 ############################ 146 def fields(self): 147 """Which fields are we interested in to produce transformed data?""" 148 return [self.course_field, self.speed_field, self.heading_field, 149 self.wind_dir_field, self.wind_speed_field] 150 151 ############################ 152 def _metadata(self): 153 """Return a dict of metadata for our derived fields.""" 154 155 metadata_fields = { 156 self.true_dir_name: { 157 'description': 'Derived true wind direction from %s, %s, %s, %s, %s' 158 % (self.course_field, self.speed_field, self.heading_field, 159 self.wind_dir_field, self.wind_speed_field), 160 'units': 'degrees', 161 'device': 'TrueWindTransform', 162 'device_type': 'DerivedTrueWindTransform', 163 'device_type_field': self.true_dir_name 164 }, 165 self.true_speed_name: { 166 'description': 'Derived true wind speed from %s, %s, %s, %s, %s' 167 % (self.course_field, self.speed_field, self.heading_field, 168 self.wind_dir_field, self.wind_speed_field), 169 'units': 'depends on conversion used for %s (%g) and %s (%g)' 170 % (self.speed_field, self.convert_speed_factor, 171 self.wind_speed_field, self.convert_wind_factor), 172 'device': 'TrueWindTransform', 173 'device_type': 'DerivedTrueWindTransform', 174 'device_type_field': self.true_speed_name 175 }, 176 self.apparent_dir_name: { 177 'description': 'Derived apparent wind speed from %s, %s, %s, %s, %s' 178 % (self.course_field, self.speed_field, self.heading_field, 179 self.wind_dir_field, self.wind_speed_field), 180 'units': 'degrees', 181 'device': 'TrueWindTransform', 182 'device_type': 'DerivedTrueWindTransform', 183 'device_type_field': self.apparent_dir_name 184 } 185 } 186 return metadata_fields 187 188 ############################ 189 def transform(self, record: Union[dict, DASRecord]): 190 """Incorporate any useable fields in this record, and if it gives 191 us a new true wind value, return the results.""" 192 193 # See if it's something we can process, and if not, try digesting 194 if not self.can_process_record(record): # inherited from BaseModule() 195 return self.digest_record(record) # inherited from BaseModule() 196 197 results = [] 198 for das_record in to_das_record_list(record): 199 # If they haven't specified specific fields we should wait for 200 # before updates, plan to emit an update after every new record 201 # we process. Otherwise, assume we're not going to update unless 202 # we see one of the named fields. 203 if not self.update_on_fields: 204 update = True 205 else: 206 update = False 207 208 timestamp = das_record.timestamp 209 if not timestamp: 210 logging.info('DASRecord is missing timestamp - skipping') 211 continue 212 213 # Get latest values for any of our fields 214 fields = das_record.fields 215 if self.course_field in fields: 216 if timestamp >= self.course_val_time: 217 self.course_val = fields.get(self.course_field) 218 self.course_val_time = timestamp 219 if self.course_field in self.update_on_fields: 220 update = True 221 222 if self.speed_field in fields: 223 if timestamp >= self.speed_val_time: 224 self.speed_val = fields.get(self.speed_field) 225 self.speed_val *= self.convert_speed_factor 226 self.speed_val_time = timestamp 227 if self.speed_field in self.update_on_fields: 228 update = True 229 230 if self.heading_field in fields: 231 if timestamp >= self.heading_val_time: 232 self.heading_val = fields.get(self.heading_field) 233 self.heading_val_time = timestamp 234 if self.heading_field in self.update_on_fields: 235 update = True 236 237 if self.wind_dir_field in fields: 238 if timestamp >= self.wind_dir_val_time: 239 self.wind_dir_val = fields.get(self.wind_dir_field) 240 self.wind_dir_val_time = timestamp 241 if self.wind_dir_field in self.update_on_fields: 242 update = True 243 244 if self.wind_speed_field in fields: 245 if timestamp >= self.wind_speed_val_time: 246 self.wind_speed_val = fields.get(self.wind_speed_field) 247 self.wind_speed_val *= self.convert_wind_factor 248 self.wind_speed_val_time = timestamp 249 if self.wind_speed_field in self.update_on_fields: 250 update = True 251 252 # Check if needed all values are present, and none are too old to use 253 if self._values_too_old(timestamp): 254 continue 255 256 # If we've not seen anything that updates fields that would 257 # trigger a new true winds value, skip rest of computation. 258 if not update: 259 logging.debug('No update needed') 260 continue 261 262 logging.debug('Computing new true winds') 263 (true_dir, true_speed, apparent_dir) = truew(crse=self.course_val, 264 cspd=self.speed_val, 265 hd=self.heading_val, 266 wdir=self.wind_dir_val, 267 zlr=self.zero_line_reference, 268 wspd=self.wind_speed_val) 269 270 logging.debug('Got true winds: dir: %s, speed: %s, apparent_dir: %s', 271 true_dir, true_speed, apparent_dir) 272 if None in (true_dir, true_speed, apparent_dir): 273 logging.info('Got invalid true winds') 274 continue 275 276 # If here, we've got a valid new true wind result 277 true_wind_fields = {self.true_dir_name: true_dir, 278 self.true_speed_name: true_speed, 279 self.apparent_dir_name: apparent_dir} 280 281 # Add in metadata if so specified and it's been long enough since 282 # we last sent it. 283 now = time.time() 284 if self.metadata_interval and \ 285 now - self.metadata_interval > self.last_metadata_send: 286 metadata = {'fields': self._metadata()} 287 self.last_metadata_send = now 288 logging.debug('Emitting metadata: %s', pformat(metadata)) 289 else: 290 metadata = None 291 292 results.append(DASRecord(timestamp=timestamp, fields=true_wind_fields, 293 metadata=metadata, data_id=self.data_id)) 294 295 logging.debug('Sending %d true wind results.', len(results)) 296 return results 297 298 ############################ 299 def _values_too_old(self, timestamp): 300 """Return true if any values are missing or too old to use.""" 301 302 if None in (self.course_val, self.speed_val, self.heading_val, 303 self.wind_dir_val, self.wind_speed_val): 304 logging.debug('Not all required values for true winds are present: ' 305 'time: %s: %s: %s, %s: %s, %s: %s, %s: %s, %s: %s', 306 timestamp, 307 self.course_field, self.course_val, 308 self.speed_field, self.speed_val, 309 self.heading_field, self.heading_val, 310 self.wind_dir_field, self.wind_dir_val, 311 self.wind_speed_field, self.wind_speed_val) 312 return True 313 314 course_max_age = self.max_field_age.get(self.course_field) 315 if (course_max_age and timestamp - self.course_val_time > course_max_age): 316 logging.debug('course_field too old - max age %g, age %g', 317 course_max_age, timestamp - self.course_val_time) 318 return True 319 320 speed_max_age = self.max_field_age.get(self.speed_field) 321 if speed_max_age: 322 if timestamp - self.speed_val_time > speed_max_age: 323 logging.debug('speed_field too old - max age %g, age %g', 324 speed_max_age, timestamp - self.speed_val_time) 325 return True 326 327 heading_max_age = self.max_field_age.get(self.heading_field) 328 if heading_max_age: 329 if timestamp - self.heading_val_time > heading_max_age: 330 logging.debug('heading_field too old - max age %g, age %g', 331 heading_max_age, timestamp - self.heading_val_time) 332 return True 333 334 wind_dir_max_age = self.max_field_age.get(self.wind_dir_field) 335 if wind_dir_max_age: 336 if timestamp - self.wind_dir_val_time > wind_dir_max_age: 337 logging.debug('wind_dir_field too old - max age %g, age %g', 338 wind_dir_max_age, timestamp - self.wind_dir_val_time) 339 return True 340 341 wind_speed_max_age = self.max_field_age.get(self.wind_speed_field) 342 if wind_speed_max_age: 343 if timestamp - self.wind_speed_val_time > wind_speed_max_age: 344 logging.debug('wind_speed_field too old - max age %g, age %g', 345 wind_speed_max_age, timestamp - self.wind_speed_val_time) 346 return True 347 348 # Everything is present, and nothing's too old... 349 return False
Transform that computes true winds from vessel course/speed/heading and anemometer relative wind speed/dir.
44 def __init__(self, 45 course_field, speed_field, heading_field, 46 wind_dir_field, wind_speed_field, 47 true_dir_name, 48 true_speed_name, 49 apparent_dir_name, 50 update_on_fields=[], 51 max_field_age={}, 52 zero_line_reference=0, 53 convert_wind_factor=1, 54 convert_speed_factor=1, 55 data_id=None, 56 metadata_interval=None, 57 **kwargs): 58 """ 59 ``` 60 course_field 61 speed_field 62 heading_field 63 wind_dir_field 64 wind_speed_field 65 Field names from which we should take values for 66 course, speed over ground, heading, relative wind speed 67 and relative wind direction. 68 69 true_dir_name 70 true_speed_name 71 apparent_dir_name 72 Names that should be given to transform output values. 73 74 update_on_fields 75 If non-empty, a list of fields, any of whose arrival should 76 trigger an output record. If None, generate output when any 77 field is updated. 78 79 max_field_age 80 If non-empty, a dict of field_name:seconds, specifying that 81 no output is to be produces if the age of any of the specified 82 names is older than the specified number of seconds. 83 84 zero_line_reference 85 Angle between bow and zero line on anemometer, referenced 86 to ship. 87 88 convert_wind_factor 89 convert_speed_factor 90 Wind speed and vessel speed may be in different units; e.g 91 wind speed in meters/sec and vessel speed in knots. Multiply 92 the respective values we get by the respective factors. 93 Typically, only one of these will be not equal to 1; e.g. to 94 output true winds as meters/sec, we'll leave convert_wind_factor 95 as 1 and specify convert_speed_factor=0.5144 96 97 metadata_interval - how many seconds between when we attach field metadata 98 to a record we send out. 99 100 data_id Optional name that will be attached to the resulting DASRecord 101 ``` 102 """ 103 super().__init__(**kwargs) # processes 'quiet' and type hints 104 105 self.course_field = course_field 106 self.speed_field = speed_field 107 self.heading_field = heading_field 108 self.wind_dir_field = wind_dir_field 109 self.wind_speed_field = wind_speed_field 110 111 self.true_dir_name = true_dir_name 112 self.true_speed_name = true_speed_name 113 self.apparent_dir_name = apparent_dir_name 114 115 self.update_on_fields = update_on_fields 116 self.max_field_age = max_field_age 117 self.field_age = {} 118 119 self.zero_line_reference = zero_line_reference 120 121 self.convert_wind_factor = convert_wind_factor 122 self.convert_speed_factor = convert_speed_factor 123 124 self.metadata_interval = metadata_interval 125 self.last_metadata_send = 0 126 127 self.data_id = data_id 128 129 # TODO: It may make sense for us to cache most recent values so 130 # that, for example, we can take single DASRecords in the 131 # transform() method and use the most recent values we've seen 132 # from previous calls. 133 self.course_val = None 134 self.speed_val = None 135 self.heading_val = None 136 self.wind_dir_val = None 137 self.wind_speed_val = None 138 139 self.course_val_time = 0 140 self.speed_val_time = 0 141 self.heading_val_time = 0 142 self.wind_dir_val_time = 0 143 self.wind_speed_val_time = 0
course_field
speed_field
heading_field
wind_dir_field
wind_speed_field
Field names from which we should take values for
course, speed over ground, heading, relative wind speed
and relative wind direction.
true_dir_name
true_speed_name
apparent_dir_name
Names that should be given to transform output values.
update_on_fields
If non-empty, a list of fields, any of whose arrival should
trigger an output record. If None, generate output when any
field is updated.
max_field_age
If non-empty, a dict of field_name:seconds, specifying that
no output is to be produces if the age of any of the specified
names is older than the specified number of seconds.
zero_line_reference
Angle between bow and zero line on anemometer, referenced
to ship.
convert_wind_factor
convert_speed_factor
Wind speed and vessel speed may be in different units; e.g
wind speed in meters/sec and vessel speed in knots. Multiply
the respective values we get by the respective factors.
Typically, only one of these will be not equal to 1; e.g. to
output true winds as meters/sec, we'll leave convert_wind_factor
as 1 and specify convert_speed_factor=0.5144
metadata_interval - how many seconds between when we attach field metadata
to a record we send out.
data_id Optional name that will be attached to the resulting DASRecord
146 def fields(self): 147 """Which fields are we interested in to produce transformed data?""" 148 return [self.course_field, self.speed_field, self.heading_field, 149 self.wind_dir_field, self.wind_speed_field]
Which fields are we interested in to produce transformed data?
189 def transform(self, record: Union[dict, DASRecord]): 190 """Incorporate any useable fields in this record, and if it gives 191 us a new true wind value, return the results.""" 192 193 # See if it's something we can process, and if not, try digesting 194 if not self.can_process_record(record): # inherited from BaseModule() 195 return self.digest_record(record) # inherited from BaseModule() 196 197 results = [] 198 for das_record in to_das_record_list(record): 199 # If they haven't specified specific fields we should wait for 200 # before updates, plan to emit an update after every new record 201 # we process. Otherwise, assume we're not going to update unless 202 # we see one of the named fields. 203 if not self.update_on_fields: 204 update = True 205 else: 206 update = False 207 208 timestamp = das_record.timestamp 209 if not timestamp: 210 logging.info('DASRecord is missing timestamp - skipping') 211 continue 212 213 # Get latest values for any of our fields 214 fields = das_record.fields 215 if self.course_field in fields: 216 if timestamp >= self.course_val_time: 217 self.course_val = fields.get(self.course_field) 218 self.course_val_time = timestamp 219 if self.course_field in self.update_on_fields: 220 update = True 221 222 if self.speed_field in fields: 223 if timestamp >= self.speed_val_time: 224 self.speed_val = fields.get(self.speed_field) 225 self.speed_val *= self.convert_speed_factor 226 self.speed_val_time = timestamp 227 if self.speed_field in self.update_on_fields: 228 update = True 229 230 if self.heading_field in fields: 231 if timestamp >= self.heading_val_time: 232 self.heading_val = fields.get(self.heading_field) 233 self.heading_val_time = timestamp 234 if self.heading_field in self.update_on_fields: 235 update = True 236 237 if self.wind_dir_field in fields: 238 if timestamp >= self.wind_dir_val_time: 239 self.wind_dir_val = fields.get(self.wind_dir_field) 240 self.wind_dir_val_time = timestamp 241 if self.wind_dir_field in self.update_on_fields: 242 update = True 243 244 if self.wind_speed_field in fields: 245 if timestamp >= self.wind_speed_val_time: 246 self.wind_speed_val = fields.get(self.wind_speed_field) 247 self.wind_speed_val *= self.convert_wind_factor 248 self.wind_speed_val_time = timestamp 249 if self.wind_speed_field in self.update_on_fields: 250 update = True 251 252 # Check if needed all values are present, and none are too old to use 253 if self._values_too_old(timestamp): 254 continue 255 256 # If we've not seen anything that updates fields that would 257 # trigger a new true winds value, skip rest of computation. 258 if not update: 259 logging.debug('No update needed') 260 continue 261 262 logging.debug('Computing new true winds') 263 (true_dir, true_speed, apparent_dir) = truew(crse=self.course_val, 264 cspd=self.speed_val, 265 hd=self.heading_val, 266 wdir=self.wind_dir_val, 267 zlr=self.zero_line_reference, 268 wspd=self.wind_speed_val) 269 270 logging.debug('Got true winds: dir: %s, speed: %s, apparent_dir: %s', 271 true_dir, true_speed, apparent_dir) 272 if None in (true_dir, true_speed, apparent_dir): 273 logging.info('Got invalid true winds') 274 continue 275 276 # If here, we've got a valid new true wind result 277 true_wind_fields = {self.true_dir_name: true_dir, 278 self.true_speed_name: true_speed, 279 self.apparent_dir_name: apparent_dir} 280 281 # Add in metadata if so specified and it's been long enough since 282 # we last sent it. 283 now = time.time() 284 if self.metadata_interval and \ 285 now - self.metadata_interval > self.last_metadata_send: 286 metadata = {'fields': self._metadata()} 287 self.last_metadata_send = now 288 logging.debug('Emitting metadata: %s', pformat(metadata)) 289 else: 290 metadata = None 291 292 results.append(DASRecord(timestamp=timestamp, fields=true_wind_fields, 293 metadata=metadata, data_id=self.data_id)) 294 295 logging.debug('Sending %d true wind results.', len(results)) 296 return results
Incorporate any useable fields in this record, and if it gives us a new true wind value, return the results.