Kinesis Agent сообщает об ошибке из данных потока Python

Я пытаюсь создать простой поток данных из кода Python, который добавляется в файл журнала, и загружать потоковые данные в кластер Redshift с помощью Kinesis Agent->Kinesis Firehose.

Я подтверждаю, что код Python работает нормально, создает потоковые данные, которые добавляются в файл журнала. Мой агент Kinesis успешно запущен, и я проверил его, а также выполнил соответствующие настройки для извлечения данных из файла журнала и передачи их в поток пожарной сигнализации.

Файл python представляет собой простой код, который случайным образом регистрирует широтные / длинные позиции. Это только для того, чтобы сделать пример потокового пожарного шланга, и, следовательно, использовать образцы данных:

latitude = 19.99
longitude = 73.78
file_n = '/tmp/random_lat_lon.log'

def generate_random_data(lat, lon, num_rows, file_name):
    with open(file_name, 'w+', 1) as output:
#        for _ in xrange(num_rows):
         while True:
            hex1 = '%012x' % random.randrange(16**12)
            flt = float(random.randint(0,100))
            dec_lat = random.random()/100
            dec_lon = random.random()/100
            output.write('%s,%.1f,%.6f,%.6f \n' % (hex1.lower(), flt, lon+dec_lon, lat+dec_lat))
            time.sleep(5)

generate_random_data(latitude, longitude, 5, file_n)

Выведите в файл random_lat_lon.log:

> 83d6c9f7a0be,25.0,73.782042,19.997504
> 18b69c5c5248,25.0,73.788921,19.995153
> 6a0d182996f0,91.0,73.783399,19.998097
> 431ba9e4f38e,0.0,73.781139,19.995481

Когда я проверяю kinesis-Agent, я вижу, что он не работает, и я получаю следующую трассировку ошибки:

(FileTailer[kinesis:python-stream:/tmp/random_lat_lon.log*]) com.amazon.kinesis.streaming.agent.tailing.FileTailer [ERROR] FileTailer[kinesis:python-stream:/tmp/random_lat_lon.log*]: Error when processing current input file or when tracking its status.
java.lang.IllegalStateException
        at com.google.common.base.Preconditions.checkState(Preconditions.java:158)
        at com.amazon.kinesis.streaming.agent.tailing.TrackedFileRotationAnalyzer.findCurrentOpenFileAfterTruncate(Unknown Source)
        at com.amazon.kinesis.streaming.agent.tailing.SourceFileTracker.updateCurrentFile(Unknown Source)
        at com.amazon.kinesis.streaming.agent.tailing.SourceFileTracker.refresh(Unknown Source)
        at com.amazon.kinesis.streaming.agent.tailing.FileTailer.updateRecordParser(Unknown Source)
        at com.amazon.kinesis.streaming.agent.tailing.FileTailer.processRecords(Unknown Source)
        at com.amazon.kinesis.streaming.agent.tailing.FileTailer.runOnce(Unknown Source)
        at com.amazon.kinesis.streaming.agent.tailing.FileTailer.run(Unknown Source)
        at com.google.common.util.concurrent.AbstractExecutionThreadService$1$2.run(AbstractExecutionThreadService.java:60)
        at com.google.common.util.concurrent.Callables$3.run(Callables.java:95)
        at java.lang.Thread.run(Thread.java:748)

Моя конфигурация kinesis-Agent.json выглядит следующим образом:

{
  "cloudwatch.emitMetrics": true,
  "kinesis.endpoint": "https://kinesis.us-east-1.amazonaws.com",
  "firehose.endpoint": "",

  "flows": [
    {
      "filePattern": "/tmp/random_lat_lon.log*",
      "kinesisStream": "python-stream"
    }
  ]
}

Это мой первый опыт лабораторной работы с Kinesis Firehose (с использованием Python). Я скучаю по чему-то, что я не мог понять.

Может кто-нибудь, пожалуйста, помогите с предложениями. Дайте мне знать, если какие-либо детали требуются.

С уважением

0 ответов

Я сталкиваюсь с той же проблемой. Выглядит как открытый выпуск от amazon-kinesis-agent. Я решил проблему, имея файл с расширением.dummy, чтобы у агента был файл для чтения в любой момент времени.