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

Environment(project_id: str, resource_data: dict)
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-(.*)')
version_pattern
worker_cpu: float
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    )
worker_memory_gb: float
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    )
worker_max_count: int
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    )
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()))
parallelism: float
76  @property
77  def parallelism(self) -> float:
78    return float(self.airflow_config_overrides.get('core-parallelism', 'inf'))
composer_version: str
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)
airflow_version: str
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)
is_composer2: bool
92  @property
93  def is_composer2(self) -> bool:
94    return self.composer_version.startswith('2')
full_path: str
96  @property
97  def full_path(self) -> str:
98    return self._resource_data['name']

Returns the full path of this resource.

Example: 'projects/gcpdiag-gke-1-9b90/zones/europe-west4-a/clusters/gke1'

state: str
100  @property
101  def state(self) -> str:
102    return self._resource_data['state']
image_version: str
104  @property
105  def image_version(self) -> str:
106    return self._resource_data['config']['softwareConfig']['imageVersion']
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'

airflow_config_overrides: dict
112  @property
113  def airflow_config_overrides(self) -> dict:
114    return self._resource_data['config']['softwareConfig'].get('airflowConfigOverrides', {})
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
def parse_full_path(self) -> Tuple[str, str]:
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)
def is_private_ip(self) -> bool:
135  def is_private_ip(self) -> bool:
136    return self._resource_data['config']['privateEnvironmentConfig'].get(
137      'enablePrivateEnvironment', False
138    )
gke_cluster: str
140  @property
141  def gke_cluster(self) -> str:
142    return self._resource_data['config']['gkeCluster']
num_schedulers: int
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    )
scheduler_cpu: float
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    )
scheduler_memory_gb: float
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    )
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