【问题标题】:Cannot delete AWS SQS message from Spark application running in EMR无法从 EMR 中运行的 Spark 应用程序中删除 AWS SQS 消息
【发布时间】:2016-06-29 12:32:22
【问题描述】:

我正在 AWS EMR 集群中运行 Apache Spark 应用程序。应用程序从 AWS SQS 检索消息,根据消息数据进行一些计算,然后删除每条消息。

我正在使用 NAT 实例的私有子网上的 VPC 中运行 EMR 集群。

我面临的问题是我无法删除该消息。我可以检索所有消息并且可以发送消息,但我无法删除它们。

我在 EMR 集群上使用以下安全性 EC2 instance profile:EMR_EC2_DefaultRole EMR role:EMR_DefaultRole

每个角色都附加了以下政策: AmazonSQSFullAccessAmazonElastiCacheFullAccessAmazonElasticMapReduceFullAccessAmazonVPCFullAccess

我认为问题出在权限上,但 AmazonSQSFullAccess 授予完全权限,所以我没有选择。

这是删除消息的 Java 代码:

public class SQSMessageBroker
{
    private AmazonSQS _amazonSqs;

    public SQSMessageBroker()
    {
        // Create the SQS client
        createSQSClient();
    }

    public void deleteMessage(String queueUrl, String receiptHandle)
        {
            DeleteMessageRequest deleteMessageRequest = new DeleteMessageRequest(queueUrl, receiptHandle);

            _amazonSqs.deleteMessage(deleteMessageRequest);
        }

  private void createSQSClient()
    {
        _amazonSqs = new AmazonSQSClient();
        _amazonSqs.setRegion(Region.getRegion(Regions.EU_WEST_1));
    }
}

SQSMessageBroker 在我的应用程序中是一个单例。 当我在本地运行相同的代码时,一切都很好。我在本地创建了一个 AWS 用户,并将密钥和密钥添加到 .aws 文件中。


编辑

经过大量研究和测试,这是我发现的:

  1. 看来这不是权限问题(至少对于由 EMR 启动的 EC2 实例而言不是)。我连接到实例,安装了 aws cli,检索到一条消息并成功删除它。
  2. _amazonSqs.deleteMessage(deleteMessageRequest); 代码不会引发任何异常。看起来好像请求超时,但没有抛出超时异常。 deleteMessage 之后的任何代码都不会执行。
  3. 我在一个单独的线程中处理每条消息,所以我为每个线程添加了一个Thread.UncaughtExceptionHandler,但那里也没有抛出异常。
  4. 我怀疑问题可能出在ReceiptHandle,更准确地说,因为我在几台机器上运行了一个Spark集群,所以我认为机器IP、名称或类似的东西被编码在ReceiptHandle中并且deleteMessage 可能是从另一台机器上执行的,所以这就是它不起作用的原因。这就是我创建一个只有一台机器的 Spark 集群的原因。遗憾的是,我仍然无法删除该消息。

【问题讨论】:

  • 是的,我有红色的。现在我开始想知道收据句柄需要多长时间才能过期,如果这可能是这种情况? Spark 应用程序需要一些时间(20-30 秒)来完成计算。
  • 遗憾的是,过期的令牌可能并非如此。我尝试从另一个队列中检索消息并立即删除它们。那没起效。我怀疑问题可能在于一些奇怪的权限问题,但我仍然不知道它是什么。
  • 将此“EU_WEST_1”更改为“EU-WEST-1”
  • 它是一个静态的 Regions 属性。它来自 AWS SKD,它适用于检索和推送消息,如果我在我的机器上运行应用程序,我也可以删除消息,所以我不认为这是问题所在。不过我会检查Regions 类中是否有EU-WEST-1 属性。

标签: java amazon-web-services amazon-ec2 amazon-sqs


【解决方案1】:

经过大量调试和测试后,我终于设法找出问题所在。

正如预期的那样,这不是权限问题。问题是,由 EMR 启动并在其上运行 Spark 应用程序的 EC2 实例包含所有 AWS 包的特定版本(包括 SQS 包)。包含包的路径被添加到 Hadoop、Yarn 和 Spark。所以当我的应用程序启动时,它使用了机器上已经存在的包,我收到了一个错误。 (错误记录在 Yarn 日志中。我花了一些时间才弄清楚。)

我正在使用 maven shade 插件为我的应用程序构建 uber jar,所以我认为我可以尝试对 AWS 包进行遮蔽(重新定位)。这将允许我将依赖项封装在我的应用程序中。可悲的是,这不起作用。亚马逊似乎在包内使用反射,并且他们对某些类的名称进行了硬编码,从而使阴影变得无用。(在我的阴影包中找不到硬编码的类)

所以经过一番挫折后,我找到了以下解决方案:

  1. 创建一个 EMR 步骤,将我的 uber jar 从 S3 下载到机器。
  2. 使用以下 spark-submit 选项创建 Spark 应用程序步骤:

--driver-class-path /path_to_your_jar/myapp.jar --class com.myapp.startapp

这里的关键是--driver-class-path 选项。你可以阅读更多关于它的信息here。基本上我将我的 uber jar 添加到 Spark 驱动程序类路径中,允许应用程序使用我的依赖项。

到目前为止,这是我找到的唯一可接受的解决方案。如果您知道另一个或更好的,请写评论或回答。

我希望这个答案对一些不幸的人有用。它会为我节省好几天的痛苦。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-09-28
    • 2020-03-24
    • 2020-12-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-04-29
    相关资源
    最近更新 更多