Skip to content
This repository was archived by the owner on Sep 19, 2024. It is now read-only.
This repository was archived by the owner on Sep 19, 2024. It is now read-only.

Problem with fs.appendFile with Statd Client ? #70

@fsck-mount

Description

@fsck-mount

This one looks so strange to me. I'm assuming this problem is because of StatD because once I remove it from code, it is working super fine. Here is the details about the issue:

I've Kafka Data floating here and there. I've node js client which consumes the data from Kafka and store it locally for 100 ms. Then it flush to disk at every 100 ms.

The code will be something similar like this:

const MainApp = require('./MainApp');

const fs = require('fs');
const path = require('path');
const uuid = require('uuid');


const zkUrl = process.env.ZK_URL;
const topics = process.env.SYNC_TOPICS_TO_CONSUME.split(' ');

const consumerOptions = {
  zookeeperUrl: zkUrl,
  groupId: process.env.KAFKA_SYNCER_GROUP_ID,
  serverPort: 3043,
  threadCount: parseInt(process.env.KAFKA_SYNCER_WORKER_COUNT || 1, 10),
  properties: { 'rebalance.max.retries': '3' }
};

function mkdirp(_filepath) {
  const dirname = path.dirname(_filepath);
  if (!fs.existsSync(dirname)) {
    mkdirp(dirname);
  }
  fs.mkdirSync(_filepath);
}


class Syncer extends MainApp {
  constructor(options) {
    if (!options) optiosn = {};
    super();
    Object.defineProperties(this, {
      bigData: {
        value: {},
        writable: true,
        enumerable: true
      }
    });

    if (options.autoSyncDb) {
      const self = this;
      const interval = parseInt(options.autoSyncDbInterval, 10) || 100;
      setTimeout(function _orbSyncDb() {
        self.syncDb();
        setTimeout(_orbSyncDb, interval);
      });
    }
  }

  pushMessage(message, partition) {
    if (this.bigData[partition] === undefined) {
      this.bigData[partition] = [message];
    } else {
      this.bigData[partition].push(message);
    }
  }

  syncDb() {
    this._bigDataCopy = this.bigData;
    this.bigData = {};

    Object.keys(this._bigDataCopy).forEach((_partition) => {
      const days = this._bigDataCopy[_partition];

      this.fileName = 'events.json';
      this.fileMessages = {};

      days.forEach((day) => {
        
        fs.readdirSync(filePath).forEach((file) => {
          if (file.indexOf('events.json') > 0) {
            const fileStats = fs.statSync(filePath + file);
            if (!fullfilePath && ((fileStats.size / 1000000) < 500)) {
              fullfilePath = filePath + file;
            }
          }
        });

        if (!fullfilePath) {
          fullfilePath = filePath + uuid() + this.fileName;
        }

        fs.appendFile(
          fullfilePath,
          day.map((_msg) => { // eslint-disable-line
            return JSON.stringify(_msg);
          }).join('\n') + '\n',
          (err) => {
            if (err) MainApp.log.error(err, [filePath, this.fileName]);
            console.log(
              "Flushed to disk: " + (this.fileMessages[msgDay] || []).length
              + " Messages,");
          }
        );
      });
    });
  }

  forceFlush() {
    this.syncDb();
  }
}

const syncer = new Syncer({
  autoSyncSchema: true,
  autoSyncSchemaInterval: 30000,
  autoSyncDb: true,
  autoSyncDbInterval: 100
});

class MessageParser extends MainApp {
  constructor(message, topic) {
    super();
    this.topic = topic;
    if (typeof message === 'string') {
      try {
        this.message = JSON.parse(message);
      } catch (e) {
        this.message = message;
        this.log.error(e, [message]);
      }
    } else {
      this.message = message;
    }
  }

  getDay (){
    return this.message.timestamp.split('T')[0];
  }
}

const onMessage = function onKafkaMessage(message, topic, callback) {
  this.event = new MessageParser(message, topic);
  syncer.pushMessage(this.event.getDay(), this.event.message);
  process.nextTick(callback);
};


MainApp.consumer.consume(topics, onMessage, consumerOptions);

process.on('SIGINT', () => {
  if (MainApp.consumer.doesExist()) MainApp.consumer.stop();
  syncer.forceFlush();
  MainApp.log.info('Consumer Stopped!');
  process.exit();
});

process.on('uncaughtException', (err) => {
  MainApp.log.error(err);
  process.exit();
});

The above one is working super fine., and the tough part starts here. When I made changes and added statsD, appendSync is not working at all,

here is the code after changes

const MainApp = require('./MainApp');

const fs = require('fs');
const path = require('path');
const uuid = require('uuid');


const statsdOptions = {
  host: process.env.STATSD_HOST,
  port: process.env.STATSD_PORT,
  prefix: 'vtap.syncer.',
  globalize: true
};

Vtap.statsD(statsdOptions);

const zkUrl = process.env.ZK_URL;
const topics = process.env.SYNC_TOPICS_TO_CONSUME.split(' ');

const consumerOptions = {
  zookeeperUrl: zkUrl,
  groupId: process.env.KAFKA_SYNCER_GROUP_ID,
  serverPort: 3043,
  threadCount: parseInt(process.env.KAFKA_SYNCER_WORKER_COUNT || 1, 10),
  properties: { 'rebalance.max.retries': '3' }
};

function mkdirp(_filepath) {
  const dirname = path.dirname(_filepath);
  if (!fs.existsSync(dirname)) {
    mkdirp(dirname);
  }
  fs.mkdirSync(_filepath);
}


class Syncer extends MainApp {
  constructor(options) {
    if (!options) optiosn = {};
    super();
    Object.defineProperties(this, {
      bigData: {
        value: {},
        writable: true,
        enumerable: true
      }
    });

    if (options.autoSyncDb) {
      const self = this;
      const interval = parseInt(options.autoSyncDbInterval, 10) || 100;
      setTimeout(function _orbSyncDb() {
        self.syncDb();
        setTimeout(_orbSyncDb, interval);
      });
    }
  }

  pushMessage(message, partition) {
    if (this.bigData[partition] === undefined) {
      this.bigData[partition] = [message];
    } else {
      this.bigData[partition].push(message);
    }
  }

  syncDb() {
    this._bigDataCopy = this.bigData;
    this.bigData = {};

    Object.keys(this._bigDataCopy).forEach((_partition) => {
      const days = this._bigDataCopy[_partition];

      this.fileName = 'events.json';
      this.fileMessages = {};

      days.forEach((day) => {
        
        fs.readdirSync(filePath).forEach((file) => {
          if (file.indexOf('events.json') > 0) {
            const fileStats = fs.statSync(filePath + file);
            if (!fullfilePath && ((fileStats.size / 1000000) < 500)) {
              fullfilePath = filePath + file;
            }
          }
        });

        if (!fullfilePath) {
          fullfilePath = filePath + uuid() + this.fileName;
        }

        fs.appendFile(
          fullfilePath,
          day.map((_msg) => { // eslint-disable-line
            return JSON.stringify(_msg);
          }).join('\n') + '\n',
          (err) => {
            if (err) {
            MainApp.log.error(err, [filePath, this.fileName]);
            }
            else {
            this.flushedMsgCount = (day || []).length;
            console.log(`Flushed to disk: ${day} Messages.`);
            }
            
          }
        );
      });
    });
  }

  forceFlush() {
    this.syncDb();
  }
}

const syncer = new Syncer({
  autoSyncSchema: true,
  autoSyncSchemaInterval: 30000,
  autoSyncDb: true,
  autoSyncDbInterval: 100
});

class MessageParser extends MainApp {
  constructor(message, topic) {
    super();
    this.topic = topic;
    if (typeof message === 'string') {
      try {
        this.message = JSON.parse(message);
      } catch (e) {
        this.message = message;
        this.log.error(e, [message]);
      }
    } else {
      this.message = message;
    }
  }

  getDay (){
    return this.message.timestamp.split('T')[0];
  }
}

const onMessage = function onKafkaMessage(message, topic, callback) {
  this.event = new MessageParser(message, topic);
  syncer.pushMessage(this.event.getDay(), this.event.message);
  process.nextTick(callback);
};


MainApp.consumer.consume(topics, onMessage, consumerOptions);

process.on('SIGINT', () => {
  if (MainApp.consumer.doesExist()) MainApp.consumer.stop();
  syncer.forceFlush();
  MainApp.log.info('Consumer Stopped!');
  process.exit();
});

process.on('uncaughtException', (err) => {
  MainApp.log.error(err);
  process.exit();
});

I don't understand what causing the issue. I don't see any error message, any console message at all. To debug I added console log every where,

Till FileFullPath it is working, after that no logs are created. Any idea why it is failing to write to disk ?

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions