1

我正在尝试使用 flink-streaming 状态后端,遵循本指南:https ://ci.apache.org/projects/flink/flink-docs-master/apis/streaming/state.html ,但我收到错误:无法解析符号“ValueState”

看了一会儿,我意识到ValueState不在我的依赖项中。相反,只有OperatorState在 org.apache.flink.api.common.state (flink-core) 中。

但是,如果我查看 Github,我会在该包中看到 ValueState:https ://github.com/apache/flink/tree/master/flink-core/src/main/java/org/apache/flink/api/common /状态

我猜我要么没有正确版本的 flink 以按照指南显示的方式使用 StateBackend,要么我有正确的版本,但 ValueState 已移至另一个 maven 依赖项。

下面是我的 pom.xml:

<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
     xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>

<groupId>test</groupId>
<artifactId>flink-streaming</artifactId>
<version>1.0-SNAPSHOT</version>
<packaging>jar</packaging>

<name>flink-streaming</name>
<url>http://maven.apache.org</url>

<properties>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
    <!--<flink.version>0.10.2</flink.version>-->
    <flink.version>0.10.2</flink.version>
    <scala.version>2.11.8</scala.version>
    <scala.dependency.version>2.11</scala.dependency.version>
</properties>


<dependencies>
    <dependency>
        <groupId>org.scala-lang</groupId>
        <artifactId>scala-library</artifactId>
        <version>${scala.version}</version>
    </dependency>
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-scala_${scala.dependency.version}</artifactId>
        <version>${flink.version}</version>
    </dependency>
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-clients_${scala.dependency.version}</artifactId>
        <version>${flink.version}</version>
    </dependency>
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-streaming-scala_${scala.dependency.version}</artifactId>
        <version>${flink.version}</version>
    </dependency>
</dependencies>

<build>
    <plugins>
        <plugin>
            <artifactId>maven-compiler-plugin</artifactId>
            <version>2.3.2</version>
            <configuration>
                <source>1.7</source>
                <target>1.7</target>
            </configuration>
        </plugin>

        <plugin>
            <groupId>org.scala-tools</groupId>
            <artifactId>maven-scala-plugin</artifactId>
            <executions>
                <execution>
                    <goals>
                        <goal>compile</goal>
                        <goal>testCompile</goal>
                    </goals>
                </execution>
            </executions>
            <configuration>
                <jvmArgs>
                    <jvmArg>-Xms64m</jvmArg>
                    <jvmArg>-Xmx1024m</jvmArg>
                </jvmArgs>
            </configuration>
        </plugin>

    </plugins>
</build>

这是我的代码:

import org.apache.flink.api.common.functions.RichFlatMapFunction;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.util.Collector;


public class CountWindowAverage extends RichFlatMapFunction<Tuple2<Long,Long>, Tuple2<Long,Long>> {
    private transient ValueState<Tuple2<Long,Long>> sum;

    @Override
    public void flatMap(Tuple2<Long,Long> input, Collector<Tuple2<Long,Long>> out) throws Exception {

    }
}

非常感谢您的帮助!

劳伦特。

4

1 回答 1

0

你是对的 Flink 版本0.10.x还没有ValueState. 如果你切换到至少一个版本1.0.0,你应该没问题。

于 2016-05-08T05:02:15.447 回答