【问题标题】:AWS Lambda function handler not inserting to AthenaAWS Lambda 函数处理程序未插入 Athena
【发布时间】:2020-02-25 09:38:21
【问题描述】:

我使用 Amazon Athena 的 sn-p 示例只是为了测试插入一些数据。我不知道为什么它不起作用,并且在语句执行完成时 CloudWatch 日志不显示任何输出。即使我将其更改为简单的选择语句,我也看不到任何输出。我知道查询、数据库和表都很好,因为当我使用 Athena 查询编辑器对其进行测试时,它的执行没有问题。

module.exports.dlr = async event => {

  let awsFileCreds = {
    accessKeyId: "XXX",
    secretAccessKey: "XXX"
  };

  let creds = new AWS.Credentials(awsFileCreds);
  AWS.config.credentials = creds;

  let client = new AWS.Athena({ region: "eu-west-1" });

  let q = Queue((id, cb) => {
    startPolling(id)
      .then(data => {
        return cb(null, data);
      })
      .catch(err => {
        console.log("Failed to poll query: ", err);
        return cb(err);
      });
  }, 5);

    const sql = "INSERT INTO delivery_receipts (status, eventid, mcc, mnc, msgcount, msisdn, received, userreference) VALUES ('TestDLR', 345345, 4353, '5345435', 234, '345754', 234, '8833')"

  makeQuery(sql)
    .then(data => {
      console.log("Row Count: ", data.length);
      console.log("DATA: ", data);
    })
    .catch(e => {
      console.log("ERROR: ", e);
    });

  function makeQuery(sql) {
    return new Promise((resolve, reject) => {
      let params = {
        QueryString: sql,
        ResultConfiguration: { OutputLocation: ATHENA_OUTPUT_LOCATION },
        QueryExecutionContext: { Database: ATHENA_DB }
      };

      client.startQueryExecution(params, (err, results) => {
        if (err) return reject(err);
        q.push(results.QueryExecutionId, (err, qid) => {
          if (err) return reject(err);
          return buildResults(qid)
            .then(data => {
              return resolve(data);
            })
            .catch(err => {
              return reject(err);
            });
        });
      });
    });
  }

  function buildResults(query_id, max, page) {
    let max_num_results = max ? max : RESULT_SIZE;
    let page_token = page ? page : undefined;
    return new Promise((resolve, reject) => {
      let params = {
        QueryExecutionId: query_id,
        MaxResults: max_num_results,
        NextToken: page_token
      };

      let dataBlob = [];
      go(params);

      function go(param) {
        getResults(param)
          .then(res => {
            dataBlob = _.concat(dataBlob, res.list);
            if (res.next) {
              param.NextToken = res.next;
              return go(param);
            } else return resolve(dataBlob);
          })
          .catch(err => {
            return reject(err);
          });
      }

      function getResults() {
        return new Promise((resolve, reject) => {
          client.getQueryResults(params, (err, data) => {
            if (err) return reject(err);
            var list = [];
            let header = buildHeader(
              data.ResultSet.ResultSetMetadata.ColumnInfo
            );
            let top_row = _.map(_.head(data.ResultSet.Rows).Data, n => {
              return n.VarCharValue;
            });
            let resultSet =
              _.difference(header, top_row).length > 0
                ? data.ResultSet.Rows
                : _.drop(data.ResultSet.Rows);
            resultSet.forEach(item => {
              list.push(
                _.zipObject(
                  header,
                  _.map(item.Data, n => {
                    return n.VarCharValue;
                  })
                )
              );
            });
            return resolve({
              next: "NextToken" in data ? data.NextToken : undefined,
              list: list
            });
          });
        });
      }
    });
  }

  function startPolling(id) {
    return new Promise((resolve, reject) => {
      function poll(id) {
        client.getQueryExecution({ QueryExecutionId: id }, (err, data) => {
          if (err) return reject(err);
          if (data.QueryExecution.Status.State === "SUCCEEDED")
            return resolve(id);
          else if (
            ["FAILED", "CANCELLED"].includes(data.QueryExecution.Status.State)
          )
            return reject(
              new Error(`Query ${data.QueryExecution.Status.State}`)
            );
          else {
            setTimeout(poll, POLL_INTERVAL, id);
          }
        });
      }
      poll(id);
    });
  }

  function buildHeader(columns) {
    return _.map(columns, i => {
      return i.Name;
    });
  }

  return { message: 'Go Serverless v1.0! Your function executed successfully!', event };
};

【问题讨论】:

    标签: javascript node.js aws-lambda amazon-athena


    【解决方案1】:

    想通了。使用 athena-express 包很容易将 aws lambda 事件与 athena 一起使用。您可以像往常一样指定您的配置并查询 athena 数据库,使用的代码比 amazon athena nodejs 示例中提供的代码少得多。

    这是我用来实现结果的代码:

    "use strict";
    
    const AthenaExpress = require("athena-express"),
        aws = require("aws-sdk");
    
    const athenaExpressConfig = {
        aws,
        db: "messaging",
        getStats: true
    };
    const athenaExpress = new AthenaExpress(athenaExpressConfig);
    
    exports.handler = async event => {
        const sqlQuery = "SELECT * FROM delivery_receipts LIMIT 3";
    
        try {
            let results = await athenaExpress.query(sqlQuery);
            return results;
        } catch (error) {
            return error;
        }
    };
    

    【讨论】:

    • 请注意,return error; 对我不起作用。我必须做return error.message; 才能得到实际的错误。
    猜你喜欢
    • 1970-01-01
    • 2020-05-02
    • 2019-02-08
    • 2020-02-13
    • 2016-09-04
    • 1970-01-01
    • 2016-08-06
    • 2021-05-20
    • 1970-01-01
    相关资源
    最近更新 更多