【问题标题】:Getting a Primary Key error in Rails using Sidekiq and Sidekiq-Cron使用 Sidekiq 和 Sidekiq-Cron 在 Rails 中获取主键错误
【发布时间】:2017-07-26 23:26:26
【问题描述】:

我有一个 Rails 项目,它使用 Sidekiq 处理工作任务,使用 Sidekiq-Cron 处理调度。不过,我遇到了问题。我构建了一个控制器(如下),它处理我所有的 API 查询、数据验证,然后将数据插入数据库。所有逻辑都正常运行。

然后我撕掉了实际将 API 数据插入数据库的代码部分,并将其移至 Job 类中。这样,Controller 方法可以简单地将所有繁重的工作交给工作。当我测试它时,所有逻辑都正常运行。

最后,我创建了一个 Job,它每分钟调用 Controller 方法,进行验证检查,然后启动另一个 Job 以保存 API 数据(如有必要)。当我这样做时,逻辑的第一部分似乎工作,它插入新的事件数据,但它检查这是否是我们第一次看到特定对象的事件的逻辑似乎失败了。结果是 PG 中的主键违规。

代码如下:

控制器

require 'date'

class MonnitOpenClosedSensorsController < ApplicationController

    def holderTester()
        #MonnitschedulerJob.perform_later(nil)
    end

    # Create Sidekiq queue to process new sensor readings
    def queueNewSensorEvents(auth_token, network_id)

        m = Monnit.new("iMonnit", 1)

        # Construct the query to select the most recent communication date for each sensor in the network
        lastEventForEachSensor = MonnitOpenClosedSensor.select('"SensorID", MAX("LastCommunicationDate") as "lastCommDate"')
        lastEventForEachSensor = lastEventForEachSensor.group("SensorID")
        lastEventForEachSensor = lastEventForEachSensor.where('"CSNetID" = ?', network_id)

        todaysDate = Date.today
        sevenDaysAgo = (todaysDate - 7)

        lastEventForEachSensor.each do |event|
            # puts event["lastCommDate"]
            recentEvent = MonnitOpenClosedSensor.select('id, "SensorID", "LastCommunicationDate"')
            recentEvent = recentEvent.where('"CSNetID" = ? AND "SensorID" = ? AND "LastCommunicationDate" = ?', network_id, event["SensorID"], event["lastCommDate"])

            recentEvent.each do |recent|
                message = m.get_extended_sensor(auth_token, recent["SensorID"])
                if message["LastDataMessageMessageGUID"] != recent["id"]
                    MonnitopenclosedsensorJob.perform_later(auth_token, network_id, message["SensorID"])
                    # puts "hi inner"
                    # puts message["LastDataMessageMessageGUID"]
                    # puts recent['id']
                    # puts recent["SensorID"]
                    # puts message["SensorID"]
                    # raise message
                end
            end
        end

        # Queue up any Sensor Events for new sensors
        # This would be sensors we've never seen before, from a Postgres standpoint
        sensors = m.get_sensor_ids(auth_token)
        sensors.each do |sensor|
            sensorCheck = MonnitOpenClosedSensor.select(:SensorID)
            # sensorCheck = MonnitOpenClosedSensor.select(:SensorID)
            sensorCheck = sensorCheck.group(:SensorID)
            sensorCheck = sensorCheck.where('"CSNetID" = ? AND "SensorID" = ?', network_id, sensor)
            # sensorCheck = sensorCheck.where('id = "?"', sensor["LastDataMessageMessageGUID"])

            if sensorCheck.any? == false
                MonnitopenclosedsensorJob.perform_later(auth_token, network_id, sensor) 
            end
        end

    end

end

以上代码中断了新传感器的传感器事件。它不识别传感器已经存在,首先发出,然后不识别它试图创建的事件已经持久化到数据库中(使用 GUID 进行比较)。

持久化数据的工作

class MonnitopenclosedsensorJob < ApplicationJob
  queue_as :default

  def perform(auth_token, network_id, sensor)
    m = Monnit.new("iMonnit", 1)
    newSensor = m.get_extended_sensor(auth_token, sensor)

    sensorRecord = MonnitOpenClosedSensor.new
    sensorRecord.SensorID = newSensor['SensorID']
    sensorRecord.MonnitApplicationID = newSensor['MonnitApplicationID']
    sensorRecord.CSNetID = newSensor['CSNetID']

    lastCommunicationDatePretty = newSensor['LastCommunicationDate'].scan(/[0-9]+/)[0].to_i / 1000.0
    nextCommunicationDatePretty = newSensor['NextCommunicationDate'].scan(/[0-9]+/)[0].to_i / 1000.0
    sensorRecord.LastCommunicationDate = Time.at(lastCommunicationDatePretty)
    sensorRecord.NextCommunicationDate = Time.at(nextCommunicationDatePretty)

    sensorRecord.id = newSensor['LastDataMessageMessageGUID']
    sensorRecord.PowerSourceID = newSensor['PowerSourceID']
    sensorRecord.Status = newSensor['Status']
    sensorRecord.CanUpdate = newSensor['CanUpdate'] == "true" ? 1 : 0
    sensorRecord.ReportInterval = newSensor['ReportInterval']
    sensorRecord.MinimumThreshold = newSensor['MinimumThreshold']
    sensorRecord.MaximumThreshold = newSensor['MaximumThreshold']
    sensorRecord.Hysteresis = newSensor['Hysteresis']
    sensorRecord.Tag = newSensor['Tag']
    sensorRecord.ActiveStateInterval = newSensor['ActiveStateInterval']
    sensorRecord.CurrentReading = newSensor['CurrentReading']
    sensorRecord.BatteryLevel = newSensor['BatteryLevel']
    sensorRecord.SignalStrength = newSensor['SignalStrength']
    sensorRecord.AlertsActive = newSensor['AlertsActive']
    sensorRecord.AccountID = newSensor['AccountID']
    sensorRecord.CreatedOn = Time.now.getutc
    sensorRecord.CreatedBy = "Monnit Open Closed Sensor Job"
    sensorRecord.LastModifiedOn = Time.now.getutc
    sensorRecord.LastModifiedBy = "Monnit Open Closed Sensor Job"

    sensorRecord.save

    sensorRecord = nil
  end
end

每分钟呼叫控制器的工作

class MonnitschedulerJob < ApplicationJob
  queue_as :default

  def perform(*args)
    m = Monnit.new("iMonnit", 1)
    getImonnitUsers = ImonnitCredential.select('"auth_token", "username", "password"')
    getImonnitUsers.each do |user|
        # puts user["auth_token"]
        # puts user["username"]
        # puts user["password"]

        if user["auth_token"] != nil
            m.logon(user["auth_token"])
        else
            auth_token = m.get_auth_token(user["username"], user["password"])
            auth_token = auth_token["Result"]
        end

        network_list = m.get_network_list(auth_token)
        network_list.each do |network|
            # puts network["NetworkID"]
            MonnitOpenClosedSensorsController.new.queueNewSensorEvents(auth_token, network["NetworkID"])
        end
    end
  end
end

对帖子的长度感到抱歉。我试图尽可能多地包含有关所涉及代码的信息。

编辑

这是扩展传感器的代码以及 JSON 响应:

def get_extended_sensor(auth_token, sensor_id)
        response = self.class.get("/json/SensorGetExtended/#{auth_token}?SensorID=#{sensor_id}")

        if response['Result'] != "Invalid Authorization Token"
            response['Result']
        else
            response['Result']
        end
    end


{
    "Method": "SensorGetExtended",
    "Result": {
        "ReportInterval": 180,
        "ActiveStateInterval": 180,
        "InactivityAlert": 365,
        "MeasurementsPerTransmission": 1,
        "MinimumThreshold": 4294967295,
        "MaximumThreshold": 4294967295,
        "Hysteresis": 0,
        "Tag": "",
        "SensorID": 189092,
        "MonnitApplicationID": 9,
        "CSNetID": 24391,
        "SensorName": "Open / Closed - 189092",
        "LastCommunicationDate": "/Date(1500999632000)/",
        "NextCommunicationDate": "/Date(1501010432000)/",
        "LastDataMessageMessageGUID": "d474b3db-d843-40ba-8e0e-8c4726b61ec2",
        "PowerSourceID": 1,
        "Status": 0,
        "CanUpdate": true,
        "CurrentReading": "Open",
        "BatteryLevel": 100,
        "SignalStrength": 84,
        "AlertsActive": true,
        "CheckDigit": "QOLP",
        "AccountID": 14728
    }
}

【问题讨论】:

    标签: ruby-on-rails ruby sidekiq


    【解决方案1】:

    一些想法:

    recentEvent = MonnitOpenClosedSensor.select('id, "SensorID", "LastCommunicationDate"') - 
    

    这没有做任何排序;您假设您在此处检索的记录是最新记录。

    m = Monnit.new("iMonnit", 1)
    newSensor = m.get_extended_sensor(auth_token, sensor)
    

    没有 get_extended_sensor 的实现细节是不可能告诉你的

    sensorRecord.id = newSensor['LastDataMessageMessageGUID']
    

    正在解决。

    您很可能收到重复的消息。将输入数据用作主键几乎不是一个好主意 - 而是在您的工作中自动生成一个 GUID,将其用作主键,然后使用 LastDataMessageMessageGUID 作为关联 ID。

    【讨论】:

    • 我添加了 get_extended_sensor 的 API 代码以及 JSON 响应。
    【解决方案2】:

    所以我遇到的问题,事实证明,如下:

    1. 从 API 中提取了一个传感器事件,并在 Sidekiq 中作为辅助作业排队。
    2. 如果队列运行速度有点慢、API 速度慢或只是需要处理大量作业,则 1 分钟轮询可能会再次触发,并将相同的传感器事件拉下并排队。
    3. 随着队列的处理,传感器事件被插入到数据库中,它的 GUID 是主键
    4. 当队列继续赶上自己时,它会遇到被安排为辅助作业的同一事件。该作业随后失败。

    我对此的解决方案是将“此 SensorID 和 GUID 是否存在于数据库中”转移到实际工作中。因此,当作业运行时,它要做的第一件事就是再次检查记录是否已经存在。这意味着我要检查两次,但这种快速检查的开销很低。

    仍然存在这样的风险,即在另一个作业插入记录之前,在将记录提交到数据库之前,检查可能会发生并通过,然后它可能会失败。但是重试会捕获它,然后在第二轮检查未验证时将其清除为成功的过程。话虽如此,但是在提取 API 数据之后进行检查。因为,从理论上讲,来自 API 数据的单个记录的数据库持久性会发生得非常快(比 API 调用发生的速度快得多),它确实降低了你不得不重试任何工作的机会...... .而且我的意思是,与第二次检查失败并触发重试相比,您中奖的机会更大。

    如果其他人有更好或更干净的解决方案,请随时将其作为辅助答案!

    【讨论】:

      猜你喜欢
      • 2022-12-03
      • 1970-01-01
      • 1970-01-01
      • 2014-10-16
      • 1970-01-01
      • 2018-08-29
      • 1970-01-01
      • 2017-12-04
      • 1970-01-01
      相关资源
      最近更新 更多