gcpdiag.queries.composer
Queries related to Composer.
class
Environment(gcpdiag.models.Resource):
29class Environment(models.Resource): 30 """Represents Composer environment""" 31 32 _resource_data: dict 33 34 def __init__(self, project_id: str, resource_data: dict): 35 super().__init__(project_id) 36 self._resource_data = resource_data 37 self.region, self.name = self.parse_full_path() 38 self.version_pattern = re.compile(r'composer-(.*)-airflow-(.*)') 39 40 @property 41 def worker_cpu(self) -> float: 42 return get_path( 43 self._resource_data, 44 ('config', 'workloadsConfig', 'worker', 'cpu'), 45 default=None, 46 ) 47 48 @property 49 def worker_memory_gb(self) -> float: 50 return get_path( 51 self._resource_data, 52 ('config', 'workloadsConfig', 'worker', 'memoryGb'), 53 default=None, 54 ) 55 56 @property 57 def worker_max_count(self) -> int: 58 return get_path( 59 self._resource_data, 60 ('config', 'workloadsConfig', 'worker', 'maxCount'), 61 default=None, 62 ) 63 64 @property 65 def worker_concurrency(self) -> float: 66 def default_value(): 67 airflow_version = self.airflow_version 68 69 if version.parse(airflow_version) < version.parse('2.3.3'): 70 return 12 * self.worker_cpu 71 else: 72 return min(32, 12 * self.worker_cpu, 8 * self.worker_memory_gb) 73 74 return float(self.airflow_config_overrides.get('celery-worker_concurrency', default_value())) 75 76 @property 77 def parallelism(self) -> float: 78 return float(self.airflow_config_overrides.get('core-parallelism', 'inf')) 79 80 @property 81 def composer_version(self) -> str: 82 v = self.version_pattern.search(self.image_version) 83 assert v is not None 84 return v.group(1) 85 86 @property 87 def airflow_version(self) -> str: 88 v = self.version_pattern.search(self.image_version) 89 assert v is not None 90 return v.group(2) 91 92 @property 93 def is_composer2(self) -> bool: 94 return self.composer_version.startswith('2') 95 96 @property 97 def full_path(self) -> str: 98 return self._resource_data['name'] 99 100 @property 101 def state(self) -> str: 102 return self._resource_data['state'] 103 104 @property 105 def image_version(self) -> str: 106 return self._resource_data['config']['softwareConfig']['imageVersion'] 107 108 @property 109 def short_path(self) -> str: 110 return f'{self.project_id}/{self.region}/{self.name}' 111 112 @property 113 def airflow_config_overrides(self) -> dict: 114 return self._resource_data['config']['softwareConfig'].get('airflowConfigOverrides', {}) 115 116 @property 117 def service_account(self) -> str: 118 sa = self._resource_data['config']['nodeConfig'].get('serviceAccount') 119 if sa is None: 120 # serviceAccount is marked as optional in REST API docs 121 # using a default GCE SA as a fallback 122 project_nr = crm.get_project(self.project_id).number 123 sa = f'{project_nr}-compute@developer.gserviceaccount.com' 124 return sa 125 126 def parse_full_path(self) -> Tuple[str, str]: 127 match = re.match(r'projects/[^/]*/locations/([^/]*)/environments/([^/]*)', self.full_path) 128 if not match: 129 raise RuntimeError(f"Can't parse full_path {self.full_path}") 130 return match.group(1), match.group(2) 131 132 def __str__(self) -> str: 133 return self.short_path 134 135 def is_private_ip(self) -> bool: 136 return self._resource_data['config']['privateEnvironmentConfig'].get( 137 'enablePrivateEnvironment', False 138 ) 139 140 @property 141 def gke_cluster(self) -> str: 142 return self._resource_data['config']['gkeCluster'] 143 144 @property 145 def num_schedulers(self) -> int: 146 return get_path( 147 self._resource_data, 148 ('config', 'workloadsConfig', 'scheduler', 'count'), 149 default=1, 150 ) 151 152 @property 153 def scheduler_cpu(self) -> float: 154 return get_path( 155 self._resource_data, 156 ('config', 'workloadsConfig', 'scheduler', 'cpu'), 157 default=None, 158 ) 159 160 @property 161 def scheduler_memory_gb(self) -> float: 162 return get_path( 163 self._resource_data, 164 ('config', 'workloadsConfig', 'scheduler', 'memoryGb'), 165 default=None, 166 )
Represents Composer environment
worker_concurrency: float
64 @property 65 def worker_concurrency(self) -> float: 66 def default_value(): 67 airflow_version = self.airflow_version 68 69 if version.parse(airflow_version) < version.parse('2.3.3'): 70 return 12 * self.worker_cpu 71 else: 72 return min(32, 12 * self.worker_cpu, 8 * self.worker_memory_gb) 73 74 return float(self.airflow_config_overrides.get('celery-worker_concurrency', default_value()))
full_path: str
Returns the full path of this resource.
Example: 'projects/gcpdiag-gke-1-9b90/zones/europe-west4-a/clusters/gke1'
short_path: str
108 @property 109 def short_path(self) -> str: 110 return f'{self.project_id}/{self.region}/{self.name}'
Returns the short name for this resource.
Note that it isn't clear from this name what kind of resource it is.
Example: 'gke1'
service_account: str
116 @property 117 def service_account(self) -> str: 118 sa = self._resource_data['config']['nodeConfig'].get('serviceAccount') 119 if sa is None: 120 # serviceAccount is marked as optional in REST API docs 121 # using a default GCE SA as a fallback 122 project_nr = crm.get_project(self.project_id).number 123 sa = f'{project_nr}-compute@developer.gserviceaccount.com' 124 return sa
COMPOSER_REGIONS =
['asia-northeast2', 'us-central1', 'northamerica-northeast1', 'us-west3', 'southamerica-east1', 'us-east1', 'asia-northeast1', 'europe-west1', 'europe-west2', 'asia-northeast3', 'us-west4', 'asia-east2', 'europe-central2', 'europe-west6', 'us-west2', 'australia-southeast1', 'europe-west3', 'asia-south1', 'us-west1', 'us-east4', 'asia-southeast1']
@caching.cached_api_call
def
get_environments( context: gcpdiag.models.Context) -> Iterable[Environment]:
215@caching.cached_api_call 216def get_environments(context: models.Context) -> Iterable[Environment]: 217 environments: List[Environment] = [] 218 if not apis.is_enabled(context.project_id, 'composer'): 219 return environments 220 api = apis.get_api('composer', 'v1', context.project_id) 221 222 for env in _query_regions_envs(COMPOSER_REGIONS, api, context.project_id, context): 223 # projects/{projectId}/locations/{locationId}/environments/{environmentId}. 224 result = re.match(r'projects/[^/]+/locations/([^/]+)/environments/([^/]+)', env['name']) 225 if not result: 226 logging.error('invalid composer name: %s', env['name']) 227 continue 228 location = result.group(1) 229 labels = env.get('labels', {}) 230 name = result.group(2) 231 if not context.match_project_resource(location=location, labels=labels, resource=name): 232 continue 233 234 environments.append(Environment(context.project_id, env)) 235 return environments