Commit 92142ab1 authored by Sergio Gonzalez Diaz's avatar Sergio Gonzalez Diaz
Browse files

Fix "Connection refused" error in QuestDB

parent 54f9b4cd
Loading
Loading
Loading
Loading
+2 −1
Original line number Diff line number Diff line
@@ -9,13 +9,14 @@ Jinja2==3.0.3
ncclient==0.6.13
p4runtime==1.3.0
paramiko==2.9.2
influx-line-protocol==0.1.4
# influx-line-protocol==0.1.4
python-dateutil==2.8.2
python-json-logger==2.0.2
pytz==2021.3
redis==4.1.2
requests==2.27.1
xmltodict==0.12.0
questdb==1.0.1

# pip's dependency resolver does not take into account installed packages.
# p4runtime does not specify the version of grpcio/protobuf it needs, so it tries to install latest one
+26 −17
Original line number Diff line number Diff line
@@ -12,18 +12,16 @@
# See the License for the specific language governing permissions and
# limitations under the License.

from influx_line_protocol import Metric
import socket
from questdb.ingress import Sender, IngressError
import requests
import json
import sys
import logging
import datetime

LOGGER = logging.getLogger(__name__)

class MetricsDB():
  def __init__(self, host, ilp_port, rest_port, table):
    self.socket=socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    self.host=host
    self.ilp_port=int(ilp_port)
    self.rest_port=rest_port
@@ -31,19 +29,30 @@ class MetricsDB():
    self.create_table()

  def write_KPI(self,time,kpi_id,kpi_sample_type,device_id,endpoint_id,service_id,kpi_value):
    self.socket.connect((self.host,self.ilp_port))
    metric = Metric(self.table)
    metric.with_timestamp(time)
    metric.add_tag('kpi_id', kpi_id)
    metric.add_tag('kpi_sample_type', kpi_sample_type)
    metric.add_tag('device_id', device_id)
    metric.add_tag('endpoint_id', endpoint_id)
    metric.add_tag('service_id', service_id)
    metric.add_value('kpi_value', kpi_value)
    str_metric = str(metric)
    str_metric += "\n"
    self.socket.sendall((str_metric).encode())
    self.socket.close()
    counter=0
    number_of_retries=10
    while (counter<number_of_retries):
      try:
        with Sender(self.host, self.ilp_port) as sender:
          sender.row(
          self.table,
          symbols={
              'kpi_id': kpi_id,
              'kpi_sample_type': kpi_sample_type,
              'device_id': device_id,
              'endpoint_id': endpoint_id,
              'service_id': service_id},
          columns={
              'kpi_value': kpi_value},
          at=datetime.datetime.fromtimestamp(time))
          sender.flush()
        counter=number_of_retries
        LOGGER.info(f"KPI written")
      except IngressError as ierr:
        # LOGGER.info(ierr)
        # LOGGER.info(f"Ingress Retry number {counter}")
        counter=counter+1


  def run_query(self, sql_query):
    query_params = {'query': sql_query, 'fmt' : 'json'}