Creating Your Own Replicator

Replicator files can be found in your Python library directory, under /ardi/consolidator.

Replicator modules are named 'dist_xxx', where 'xxx' is the type name used in the configuration file. Note that originally the replication engine was called the Live Data Distributor, but due to confusion with the name of an existing ARDI addon, it was renamed. However, class names and file-prefixes continue to refer to 'dist' and 'distribution'.

Installations will include a dist_custom to use as a starting point for creating your own live data replication channel.

To prevent output methods from causing disruption to ARDI data flow, replication plugins work using a queue.

When new incoming data arrives, the plugins sendData function is called. This doesn't actually send the information, but instead adds it to an internal queue to be sent in the future.

Next, updateData will be called for each individual point of data that was set with sendData. The updateData function may actually transmit the information over-the-line.

Alternatively, the updateBatchComplete function will be run after a certain number of individual updates have been processed. If your data source is consolidating information to reduce the overall number of writes, it is better to write the information here.

The script will look like the one below…

class CustomReplicator(dist.Distribution):
  def __init__(self, core):
      super(CustomDistribution, self).__init__(core)               
  def configure(self, config):
      self.config = config
      #Read In Configuration Details Here
  def isConnected(self):
      #Return 'true' if connected
      return True
  def connect(self):
      #Retrurn 'true' if connection is successful
      return True
  def sendData(self, code, value):
      #This actually QUEUES data. If you need to re-format or manipulate data
      #  before adding it to the queue, here's where you do that.
      self.queueData([str(code),str(value)])    
  def sendUpdate(self, payload):
      #This is called when the queue processes an individual item added with 'sendData'.
      #  You can send this right now if pushing out individual pieces of information.
      #  If you're consolidating information to be pushed out in a batch, write the data
      #  into memory/storage here.
      return True
      
  def updateBatchComplete(self):
      #This is called when the queue has been processed. If you're sending batched/
      #   consolidated updates, this is the place to transmit them.
      pass

You will need to fill in the functions above.

Techniques

There are two main ways these replications will usually work - sending out every individual update or sending out batched updates.

If every temperature, pressure and on/off status is going to be its own message, then you should transmit the messages in sendUpdate, as this is called for each new value.

If you'd instead like to combine updates so that they are sent less frequently, you can store the information in memory in sendUpdate and instead send the data to the destination in the updateBatchComplete function.

Example

Here's an example of writing each item to a database server individually. Note that you wouldn't use this code in production, as it does not properly validate the incoming numbers..

    def sendUpdate(self, payload):
        cur = self.conn.cursor()
                
        cur.execute("INSERT INTO " + self.table + " (" + self.namefield + "," + self.value + "," + self.time + ") VALUES ('" + payload[0])+ "', " + payload[1] + ")")            
        cur.execute("COMMIT")
        return True

And here's a more efficient method that batches the updates in sets of 50.

    def sendUpdate(self, payload):
        self.queue.append(payload)
        if len(self.queue) > 50:
              self.updateBatchComplete()
        return True
        
    def updateBatchComplete(self):
        cur = self.conn.cursor()
                
        cur.execute("INSERT INTO " + self.table + " (" + self.namefield + "," + self.value + "," + self.time + ") VALUES " + ",".join(self.queue))            
        cur.execute("COMMIT")
        self.queue = []

Internal Functions

Internally, each of these replication methods run in their own thread and maintain their own queue of messages to be transmitted. This means that a disruption/slowdown in one destination shouldn't have effects on the performance of others.

The main potential issue is excessive queue size. If possible, steps should be taken to prevent the queue from reaching extreme size. This may involve purging the queue of older records, duplicates etc. during outages.

Each replication method has exception handling to avoid errors causing the whole Consolidator from failing. You may need to add your own exception handling to understand any errors, which may otherwise be captured and quietly discarded.