]> git-server-git.apps.pok.os.sepia.ceph.com Git - ceph.git/commitdiff
mgr, mgr/telemetry: add access to osd commands in mgr and workload metrics to telemetry
authorLaura Flores <lflores@redhat.com>
Fri, 23 Jul 2021 17:40:27 +0000 (17:40 +0000)
committerLaura Flores <lflores@redhat.com>
Wed, 28 Jul 2021 16:09:00 +0000 (16:09 +0000)
In the manager module (src/pybind/mgr/mgr_module.py), there was a command that allowed access to monitor commands, but there was nothing to access osd commands. I added a new method called osd_command that allows access to the output of commands like "ceph osd.0 perf histogram dump".

In the telemetry module (src/pybind/mgr/telemetry/module.py), I added information that may help to characterize a workload. These changes include:
- io rate: delta in pg stats since last report, including a pg stats summary and some store stats
- stats sum per pool: summed up stats from every pg in a pool
- perf histograms: histogram dump of latency and request size for each osd

Signed-off-by: Laura Flores <lflores@redhat.com>
src/pybind/mgr/mgr_module.py
src/pybind/mgr/telemetry/module.py

index f7da34b771344a15bc07b13a67be7d13a33758b4..60a405847d2b47a89c6de6cb2e23c21fc854e2c3 100644 (file)
@@ -1499,6 +1499,20 @@ class MgrModule(ceph_module.BaseMgrModule, MgrModuleLoggingMixin):
 
         return r
 
+    def osd_command(self, cmd_dict: dict, inbuf: Optional[str] = None) -> Tuple[int, str, str]:
+
+        t1 = time.time()
+        result = CommandResult()
+        self.send_command(result, "osd", cmd_dict['id'], json.dumps(cmd_dict), "", inbuf)
+        r = result.wait()
+        t2 = time.time()
+
+        self.log.debug("osd_command: '{0}' -> {1} in {2:.3f}s".format(
+            cmd_dict['prefix'], r[0], t2 - t1
+        ))
+
+        return r
+
     def send_command(
             self,
             result: CommandResult,
index 3e795bd51b1efdd097e843d75e0ed6a7cd55f97f..e40ee9083a857929d1e70c4ab4d72affa621c62b 100644 (file)
@@ -302,6 +302,92 @@ class Module(MgrModule):
             'active_changed': sorted(list(active)),
         }
 
+    def get_stat_sum_per_pool(self) -> List[dict]:
+        # Initialize 'result' list
+        result: List[dict] = []
+
+        # Create a list of tuples containing pool ids and their associated names
+        # that will later act as a queue, i.e.:
+        #   pool_queue = [('1', '.mgr'), ('2', 'cephfs.a.meta'), ('3', 'cephfs.a.data')]
+        osd_map = self.get('osd_map')
+        pool_queue: List[tuple] = []
+        for pool in osd_map['pools']:
+            pool_queue.append((str(pool['pool']), pool['pool_name']))
+
+        # Populate 'result', i.e.:
+        #   {
+        #       'pool_id': '1'
+        #       'pool_name': '.mgr'
+        #       'stats_sum': {
+        #           'num_bytes': 36,
+        #           'num_bytes_hit_set_archive': 0,
+        #           ...
+        #           'num_write_kb': 0
+        #           }
+        #       }
+        #   }
+        while pool_queue:
+            # Pop a pool out of pool_queue
+            curr_pool = pool_queue.pop(0)
+
+            # Get the current pool's id and name
+            curr_pool_id, curr_pool_name = curr_pool[0], curr_pool[1]
+
+            # Initialize a dict that will hold aggregated stats for the current pool
+            compiled_stats_dict: Dict[str, dict] = defaultdict(lambda: defaultdict(int))
+
+            # Find out which pgs belong to the current pool and add up
+            # their stats
+            pg_dump = self.get('pg_dump')
+            for pg in pg_dump['pg_stats']:
+                pool_id = pg['pgid'][0:1]
+                if pool_id == curr_pool_id:
+                    for metric in pg['stat_sum']:
+                        compiled_stats_dict['pool_id'] = int(pool_id)
+                        compiled_stats_dict['pool_name'] = curr_pool_name
+                        compiled_stats_dict['stats_sum'][metric] += pg['stat_sum'][metric]
+                else:
+                    continue
+            # 'compiled_stats_dict' now holds all stats pertaining to
+            # the current pool. Adding it to the list of results.
+            result.append(compiled_stats_dict)
+
+        return result
+
+    def get_osd_histograms(self) -> Dict[str, dict]:
+        # Initialize result dict
+        result: Dict[str, dict] = defaultdict()
+
+        # You can get a list of osd ids from the metadata
+        osd_metadata = self.get('osd_metadata')
+        
+        # Grab output from the "osd.x perf histogram dump" command
+        for osd_id in osd_metadata:
+            cmd_dict = {
+                'prefix': 'perf histogram dump',
+                'id': str(osd_id),
+                'format': 'json'
+            }
+            r, outb, outs = self.osd_command(cmd_dict)
+            # Check for invalid calls
+            if r != 0:
+                continue
+            else:
+                try:
+                    # This is where the histograms will land
+                    # if there are any
+                    dump = json.loads(outb)
+                    result['osd.'+osd_id] = dump
+                # Sometimes, json errors occur if
+                # you give it an empty string
+                except json.decoder.JSONDecodeError:
+                    continue
+
+        return result
+
+    def get_io_rate(self) -> dict:
+        return self.get('io_rate')
+
     def gather_crashinfo(self) -> List[Dict[str, str]]:
         crashlist: List[Dict[str, str]] = list()
         errno, crashids, err = self.remote('crash', 'ls')
@@ -731,6 +817,9 @@ class Module(MgrModule):
 
         if 'perf' in channels:
             report['perf_counters'] = self.gather_perf_counters()
+            report['stat_sum_per_pool'] = self.get_stat_sum_per_pool()
+            report['io_rate'] = self.get_io_rate()
+            report['osd_perf_histograms'] = self.get_osd_histograms()
 
         # NOTE: We do not include the 'device' channel in this report; it is
         # sent to a different endpoint.