【问题标题】:Does Apache Flink AWS S3 Sink require Hadoop for local testing?Apache Flink AWS S3 Sink 是否需要 Hadoop 进行本地测试?
【发布时间】:2017-05-14 06:27:03
【问题描述】:

我对 Apache Flink 比较陌生,我正在尝试创建一个简单的项目来生成 AWS S3 存储桶的文件。根据文档,我似乎需要安装 Hadoop 才能执行此操作。

如何设置本地环境以允许我测试此功能?我已经在本地安装了 Apache Flink 和 Hadoop。我已经为 Hadoop 的 core-site.xml 配置添加了必要的更改,并将我的 HADOOP_CONF 路径添加到我的 flink.yaml 配置中。当我尝试通过 Flink UI 在本地提交我的工作时,我总是收到一个错误

2016-12-29 16:03:49,861 INFO  org.apache.flink.util.NetUtils                                - Unable to allocate on port 6123, due to error: Address already in use
2016-12-29 16:03:49,862 ERROR org.apache.flink.runtime.jobmanager.JobManager                - Failed to run JobManager.
java.lang.RuntimeException: Unable to do further retries starting the actor system
    at org.apache.flink.runtime.jobmanager.JobManager$.retryOnBindException(JobManager.scala:2203)
    at org.apache.flink.runtime.jobmanager.JobManager$.runJobManager(JobManager.scala:2143)
    at org.apache.flink.runtime.jobmanager.JobManager$.main(JobManager.scala:2040)
    at org.apache.flink.runtime.jobmanager.JobManager.main(JobManager.scala)

我假设我在环境设置方面遗漏了一些东西。是否可以在本地执行此操作?任何帮助,将不胜感激。

【问题讨论】:

  • 检查端口 6123 是否正在使用。如果不是,请禁用您的防火墙/iptables。

标签: hadoop amazon-s3 apache-flink flink-streaming


【解决方案1】:

虽然您需要 Hadoop 库,但您不必安装 Hadoop 即可在本地运行并写入 S3。我只是碰巧尝试编写基于 Avro 模式的 Parquet 输出并生成 SpecificRecord 到 S3。我正在通过 SBT 和 Intellij Idea 在本地运行以下代码的一个版本。所需零件:

1) 使用以下文件指定所需的 Hadoop 属性(注意:不建议定义 AWS 访问密钥/秘密密钥。最好在具有适当 IAM 角色的 EC2 实例上运行以读取/写入您的 S3 存储桶。但需要本地进行测试)

<configuration>
    <property>
        <name>fs.s3.impl</name>
        <value>org.apache.hadoop.fs.s3a.S3AFileSystem</value>
    </property>

    <!-- Comma separated list of local directories used to buffer
         large results prior to transmitting them to S3. -->
    <property>
        <name>fs.s3a.buffer.dir</name>
        <value>/tmp</value>
    </property>

    <!-- set your AWS ID using key defined in org.apache.hadoop.fs.s3a.Constants -->
    <property>
        <name>fs.s3a.access.key</name>
        <value>YOUR_ACCESS_KEY</value>
    </property>

    <!-- set your AWS access key -->
    <property>
        <name>fs.s3a.secret.key</name>
        <value>YOUR_SECRET_KEY</value>
    </property>
</configuration>

2) 进口: 导入 com.uebercomputing.eventrecord.EventOnlyRecord

import org.apache.flink.api.scala.hadoop.mapreduce.HadoopOutputFormat
import org.apache.flink.api.scala.{ExecutionEnvironment, _}

import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat
import org.apache.hadoop.conf.{Configuration => HadoopConfiguration}
import org.apache.hadoop.fs.Path
import org.apache.hadoop.mapreduce.Job

import org.apache.parquet.avro.AvroParquetOutputFormat

3) Flink 代码使用 HadoopOutputFormat 和上面的配置:

    val events: DataSet[(Void, EventOnlyRecord)] = ...

    val hadoopConfig = getHadoopConfiguration(hadoopConfigFile)

    val outputFormat = new AvroParquetOutputFormat[EventOnlyRecord]
    val outputJob = Job.getInstance

    //Note: AvroParquetOutputFormat extends FileOutputFormat[Void,T]
    //so key is Void, value of type T - EventOnlyRecord in this case
    val hadoopOutputFormat = new HadoopOutputFormat[Void, EventOnlyRecord](
      outputFormat,
      outputJob
    )

    val outputConfig = outputJob.getConfiguration
    outputConfig.addResource(hadoopConfig)
    val outputPath = new Path("s3://<bucket>/<dir-prefix>")
    FileOutputFormat.setOutputPath(outputJob, outputPath)
    AvroParquetOutputFormat.setSchema(outputJob, EventOnlyRecord.getClassSchema)

    events.output(hadoopOutputFormat)

    env.execute

    ...

    def getHadoopConfiguration(hadoodConfigPath: String): HadoopConfiguration = {
      val hadoopConfig = new HadoopConfiguration()
      hadoopConfig.addResource(new Path(hadoodConfigPath))
      hadoopConfig
    }

4) 构建依赖和使用的版本:

    val awsSdkVersion = "1.7.4"
    val hadoopVersion = "2.7.3"
    val flinkVersion = "1.1.4"

    val flinkDependencies = Seq(
      ("org.apache.flink" %% "flink-scala" % flinkVersion),
      ("org.apache.flink" %% "flink-hadoop-compatibility" % flinkVersion)
    )

    val providedFlinkDependencies = flinkDependencies.map(_ % "provided")

    val serializationDependencies = Seq(
      ("org.apache.avro" % "avro" % "1.7.7"),
      ("org.apache.avro" % "avro-mapred" % "1.7.7").classifier("hadoop2"),
      ("org.apache.parquet" % "parquet-avro" % "1.8.1")
    )

    val s3Dependencies = Seq(
      ("com.amazonaws" % "aws-java-sdk" % awsSdkVersion),
      ("org.apache.hadoop" % "hadoop-aws" % hadoopVersion)
    )

编辑使用 writeAsText 到 S3:

1) 创建一个 Hadoop 配置目录(将其引用为 hadoop-conf-dir),其中包含一个文件 core-site.xml。

例如:

mkdir /home/<user>/hadoop-config
cd /home/<user>/hadoop-config
vi core-site.xml

#content of core-site.xml 
<configuration>
    <property>
        <name>fs.s3.impl</name>
        <value>org.apache.hadoop.fs.s3a.S3AFileSystem</value>
    </property>

    <!-- Comma separated list of local directories used to buffer
         large results prior to transmitting them to S3. -->
    <property>
        <name>fs.s3a.buffer.dir</name>
        <value>/tmp</value>
    </property>

    <!-- set your AWS ID using key defined in org.apache.hadoop.fs.s3a.Constants -->
    <property>
        <name>fs.s3a.access.key</name>
        <value>YOUR_ACCESS_KEY</value>
    </property>

    <!-- set your AWS access key -->
    <property>
        <name>fs.s3a.secret.key</name>
        <value>YOUR_SECRET_KEY</value>
    </property>
</configuration>

2) 创建一个目录(将其引用为 flink-conf-dir),其中包含一个文件 flink-conf.yaml。

例如:

mkdir /home/<user>/flink-config
cd /home/<user>/flink-config
vi flink-conf.yaml

//content of flink-conf.yaml - continuing earlier example
fs.hdfs.hadoopconf: /home/<user>/hadoop-config

3) 编辑用于运行 S3 Flink 作业的 IntelliJ Run 配置 - 运行 - 编辑配置 - 并添加以下环境变量:

FLINK_CONF_DIR and set it to your flink-conf-dir

Continuing the example above:
FLINK_CONF_DIR=/home/<user>/flink-config

4) 使用该环境变量集运行代码:

events.writeAsText("s3://<bucket>/<prefix-dir>")

env.execute

【讨论】:

  • 感谢您的回复。有没有一种方法可以将我的本地 java 执行指向 hadoop 配置文件而不定义 outputPath。根据文档,我似乎应该能够执行以下操作: messageStream.writeAsText("s3://...");但是当我通过 IntelliJ 运行本地执行时,它不知道该文件在哪里。我似乎也找不到任何允许我在运行时设置它的 flink 操作。
  • 问题是调用 writeAsText 时使用的默认 HadoopFileSystem 并不“了解”s3 文件系统。请参阅上面对我的原始答案的编辑。
  • 所以我认为我一切正常,但我的 S3 存储桶的访问出现问题。出现此错误:com.amazonaws.services.s3.model.AmazonS3Exception:状态代码:403,AWS 服务:Amazon S3,AWS 请求 ID:**********,AWS 错误代码:null,AWS 错误消息:禁止,S3 扩展请求 ID:我不确定它为什么会出现访问错误,因为我的应用程序中使用的密钥与创建 S3 存储桶的帐户相同。似乎现在一切都在 flink 方面工作。如果您对我收到此错误的原因有任何提示,请告诉我。再次感谢!
  • @medium 创建存储桶与将对象放入存储桶的 IAM 权限可能不同?见docs.aws.amazon.com/AmazonS3/latest/dev/…。您可以使用这些键从 AWS CLI (aws.amazon.com/cli) 写入吗? "aws s3 cp s3:///" 听起来 Flink 部分已经设置好了。
【解决方案2】:

我必须执行以下操作才能在本地运行我的 flink 作业,该作业下沉到 S3:

1- 将 flink-s3-fs-hadoop-1.9.1.jar 添加到我的 flink/plugins/flink-s3-fs-hadoop 目录

2- 修改 flink/conf/flink-conf.yaml 以包含 s3.access-key:AWS_ACCESS_KEY s3.secret-key:AWS_SECRET_KEY fs.hdfs.hadoopconf: /etc/hadoop-config

我在 hadoop-config 文件夹中有 core-site.xml 文件,但它不包含任何配置,因此可能不需要 fs.hdfs.hadoopconf。

【讨论】:

    【解决方案3】:

    在 sbt 中我只需要添加 S3 库依赖项就可以像本地文件系统一样使用它

    SBT 文件:

    "org.apache.flink" % "flink-s3-fs-hadoop" % flinkVersion.value

    阅读示例:

        public static void main(String[] args) throws Exception {
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
    
        DataStream<String> text = env.readTextFile("s3://etl-data-ia/test/fileStreamTest.csv");
        text.print();
        env.execute("test");}
    

    【讨论】:

      【解决方案4】:

      基于该链接https://ci.apache.org/projects/flink/flink-docs-release-1.13/docs/deployment/filesystems/s3/#hadooppresto-s3-file-systems-plugins

      要使用 flink-s3-fs-hadoop 插件,您应该复制相应的 JAR 文件从 opt 目录到 Flink 的 plugins 目录 启动 Flink 之前的分发。

      我知道的另一种方法是通过环境变量启用它 ENABLE_BUILT_IN_PLUGINS="flink-s3-fs-hadoop-[flink-version].jar"

      例如:flink-s3-fs-hadoop-1.12.2.jar

      对于这两种方式,我们都必须在 flink-conf.yaml 文件中定义 S3 配置

      Flink 会在内部将其转换回 fs.s3a.connection.maximum。无需使用 Hadoop 的 XML 配置文件传递配置参数。

      s3.endpoint: <end-point>
      s3.path.style.access : true
      

      对于 AWS 凭证,它们必须在环境变量中提供,或者。在 flink-conf.yaml 中配置

      s3.endpoint: <end-point>
      s3.path.style.access : true
      s3.access-key: <key>
      s3.secret-key: <value>
      s3.region: <region>
      

      一旦完成,您就可以从 S3 中读取 @EyalP 提到的内容,或写入 S3(即使用数据集)

      dataset.map(new MapToJsonString())
                      .writeAsText("s3://....",
                              FileSystem.WriteMode.OVERWRITE);
      

      如果您想在本地测试它(没有真实的 AWS 账户),我建议您查看localstack。它完全支持各种 AWS 服务(包括 S3)。如果您这样做,则不需要 AWS 凭证(可能提供为空),并且端点将是 localstack 本身。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2014-03-18
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多