View a markdown version of this page

Erstellen Sie eine Managed Service for Apache Flink-Anwendung und führen Sie sie aus - Managed Service für Apache Flink

Die vorliegende Übersetzung wurde maschinell erstellt. Im Falle eines Konflikts oder eines Widerspruchs zwischen dieser übersetzten Fassung und der englischen Fassung (einschließlich infolge von Verzögerungen bei der Übersetzung) ist die englische Fassung maßgeblich.

Erstellen Sie eine Managed Service for Apache Flink-Anwendung und führen Sie sie aus

In diesem Schritt erstellen Sie eine Managed Service for Apache Flink-Anwendung mit Kinesis-Datenströmen als Quelle und Senke.

Erstellen Sie abhängige Ressourcen

Bevor Sie für diese Übung eine Anwendung von Managed Service für Apache Flink erstellen, erstellen Sie die folgenden abhängigen Ressourcen:

  • Zwei Kinesis-Datenströme für Eingabe und Ausgabe

  • Ein Amazon S3-Bucket zum Speichern des Anwendungscodes

    Anmerkung

    In diesem Tutorial wird davon ausgegangen, dass Sie Ihre Anwendung in der Region us-east-1 US East (Nord-Virginia) bereitstellen. Wenn Sie eine andere Region verwenden, passen Sie alle Schritte entsprechend an.

Erstellen Sie zwei Amazon Kinesis-Datenstreams

Bevor Sie für diese Übung eine Anwendung von Managed Service für Apache Flink erstellen, erstellen Sie zwei Kinesis Data Streams (ExampleInputStream und ExampleOutputStream). Ihre Anwendung verwendet diese Streams für die Quell- und Ziel-Streams der Anwendung.

Sie können diese Streams entweder mit der Amazon Kinesis-Konsole oder mit dem folgenden AWS CLI Befehl erstellen. Anweisungen für die Konsole finden Sie unter Erstellen und Aktualisieren von Datenströmen im Amazon Kinesis Data Streams Entwicklerhandbuch. Um die Streams mithilfe von zu erstellen AWS CLI, verwenden Sie die folgenden Befehle und passen Sie sie an die Region an, die Sie für Ihre Anwendung verwenden.

Um die Datenströme zu erstellen (AWS CLI)
  1. Verwenden Sie den folgenden Amazon create-stream AWS CLI Kinesis-Befehl, um den ersten Stream (ExampleInputStream) zu erstellen:

    $ aws kinesis create-stream \ --stream-name ExampleInputStream \ --shard-count 1 \ --region us-east-1 \
  2. Um den zweiten Stream zu erstellen, den die Anwendung zum Schreiben der Ausgabe verwendet, führen Sie denselben Befehl aus und ändern Sie den Stream-Namen inExampleOutputStream:

    $ aws kinesis create-stream \ --stream-name ExampleOutputStream \ --shard-count 1 \ --region us-east-1 \

Erstellen Sie einen Amazon S3-Bucket für den Anwendungscode

Sie können ein Amazon-S3-Bucket mithilfe der Konsole erstellen. Informationen zum Erstellen eines Amazon S3-Buckets mithilfe der Konsole finden Sie unter Erstellen eines Buckets im Amazon S3-Benutzerhandbuch. Benennen Sie den Amazon S3-Bucket mit einem global eindeutigen Namen, indem Sie beispielsweise Ihren Anmeldenamen anhängen.

Anmerkung

Stellen Sie sicher, dass Sie den Bucket in der Region erstellen, die Sie für dieses Tutorial verwenden (us-east-1).

Sonstige Ressourcen

Wenn Sie Ihre Anwendung erstellen, erstellt Managed Service for Apache Flink automatisch die folgenden CloudWatch Amazon-Ressourcen, sofern sie nicht bereits vorhanden sind:

  • Eine Protokollgruppe mit dem Namen /AWS/KinesisAnalytics-java/<my-application>

  • Einen Protokollstream mit dem Namen kinesis-analytics-log-stream

Einrichten der lokalen Entwicklungsumgebung

Für Entwicklung und Debugging können Sie die Apache Flink-Anwendung auf Ihrem Computer direkt von der IDE Ihrer Wahl aus ausführen. Alle Apache Flink-Abhängigkeiten werden wie normale Java-Abhängigkeiten mit Apache Maven behandelt.

Anmerkung

Auf Ihrem Entwicklungscomputer müssen Sie Java JDK 11, Maven und Git installiert haben. Wir empfehlen, eine Entwicklungsumgebung wie Eclipse Java Neon oder IntelliJ IDEA zu verwenden. Um zu überprüfen, ob Sie alle Voraussetzungen erfüllen, finden Sie unter. Erfüllen Sie die Voraussetzungen für den Abschluss der Übungen Sie müssen keinen Apache Flink-Cluster auf Ihrem Computer installieren.

Authentifizieren Sie Ihre AWS Sitzung

Die Anwendung verwendet Kinesis-Datenstreams, um Daten zu veröffentlichen. Wenn Sie lokal ausgeführt werden, benötigen Sie eine gültige AWS authentifizierte Sitzung mit Schreibberechtigungen für den Kinesis-Datenstrom. Gehen Sie wie folgt vor, um Ihre Sitzung zu authentifizieren:

  1. Wenn Sie das AWS CLI und ein benanntes Profil mit gültigen Anmeldeinformationen nicht konfiguriert haben, finden Sie weitere Informationen unter. Richten Sie das ein AWS Command Line Interface (AWS CLI)

  2. Vergewissern Sie sich, dass Ihr System korrekt konfiguriert AWS CLI ist und Ihre Benutzer über Schreibberechtigungen für den Kinesis-Datenstream verfügen, indem Sie den folgenden Testdatensatz veröffentlichen:

    $ aws kinesis put-record --stream-name ExampleOutputStream --data TEST --partition-key TEST
  3. Wenn Ihre IDE über ein Plug-in verfügt, in das Sie integrieren können AWS, können Sie es verwenden, um die Anmeldeinformationen an die in der IDE ausgeführte Anwendung zu übergeben. Weitere Informationen finden Sie unter AWS Toolkit für IntelliJ IDEA und AWS Toolkit für Eclipse.

Laden Sie den Apache Flink Streaming-Java-Code herunter und untersuchen Sie ihn

Der Java-Anwendungscode für dieses Beispiel ist unter GitHub verfügbar. Zum Herunterladen des Anwendungscodes gehen Sie wie folgt vor:

  1. Klonen Sie das Remote-Repository, indem Sie den folgenden Befehl verwenden:

    git clone https://github.com/aws-samples/amazon-managed-service-for-apache-flink-examples.git
  2. Navigieren Sie zum amazon-managed-service-for-apache-flink-examples/tree/main/java/GettingStarted Verzeichnis .

Überprüfen Sie die Anwendungskomponenten

Die Anwendung ist vollständig in der com.amazonaws.services.msf.BasicStreamingJob Klasse implementiert. Die main() Methode definiert den Datenfluss, um die Streaming-Daten zu verarbeiten und auszuführen.

Anmerkung

Für ein optimiertes Entwicklererlebnis ist die Anwendung so konzipiert, dass sie ohne Codeänderungen sowohl auf Amazon Managed Service für Apache Flink als auch lokal für die Entwicklung in Ihrer IDE ausgeführt werden kann.

  • Um die Laufzeitkonfiguration so zu lesen, dass sie funktioniert, wenn sie in Amazon Managed Service für Apache Flink und in Ihrer IDE ausgeführt wird, erkennt die Anwendung automatisch, ob sie lokal lokal in der IDE ausgeführt wird. In diesem Fall lädt die Anwendung die Laufzeitkonfiguration anders:

    1. Wenn die Anwendung feststellt, dass sie in Ihrer IDE im Standalone-Modus ausgeführt wird, erstellen Sie die application_properties.json Datei, die im Ressourcenordner des Projekts enthalten ist. Der Inhalt der Datei folgt.

    2. Wenn die Anwendung in Amazon Managed Service für Apache Flink ausgeführt wird, lädt das Standardverhalten die Anwendungskonfiguration aus den Laufzeiteigenschaften, die Sie in der Anwendung Amazon Managed Service für Apache Flink definieren. Siehe Erstellen und konfigurieren Sie die Anwendung Managed Service für Apache Flink.

      private static Map<String, Properties> loadApplicationProperties(StreamExecutionEnvironment env) throws IOException { if (env instanceof LocalStreamEnvironment) { LOGGER.info("Loading application properties from '{}'", LOCAL_APPLICATION_PROPERTIES_RESOURCE); return KinesisAnalyticsRuntime.getApplicationProperties( BasicStreamingJob.class.getClassLoader() .getResource(LOCAL_APPLICATION_PROPERTIES_RESOURCE).getPath()); } else { LOGGER.info("Loading application properties from Amazon Managed Service for Apache Flink"); return KinesisAnalyticsRuntime.getApplicationProperties(); } }
  • Die main() Methode definiert den Anwendungsdatenfluss und führt ihn aus.

    • Initialisiert die Standard-Streaming-Umgebungen. In diesem Beispiel zeigen wir, wie Sie sowohl die API für die StreamExecutionEnvironment Verwendung mit der API als auch die DataSteam API für die Verwendung mit SQL und der StreamTableEnvironment Tabellen-API erstellen. Die beiden Umgebungsobjekte sind zwei separate Verweise auf dieselbe Laufzeitumgebung, um unterschiedliche APIs zu verwenden.

      StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
    • Laden Sie die Konfigurationsparameter der Anwendung. Dadurch werden sie automatisch von der richtigen Stelle geladen, je nachdem, wo die Anwendung ausgeführt wird:

      Map<String, Properties> applicationParameters = loadApplicationProperties(env);
    • Die Anwendung definiert eine Quelle, die den Kinesis Consumer Connector verwendet, um Daten aus dem Eingabestream zu lesen. Die Konfiguration des Eingabestreams ist im PropertyGroupId = InputStream0 definiert. Der Name und die Region des Streams sind in den aws.region jeweils angegebenen stream.name Eigenschaften enthalten. Der Einfachheit halber liest diese Quelle die Datensätze als Zeichenfolge.

      private static FlinkKinesisConsumer<String> createSource(Properties inputProperties) { String inputStreamName = inputProperties.getProperty("stream.name"); return new FlinkKinesisConsumer<>(inputStreamName, new SimpleStringSchema(), inputProperties); } ... public static void main(String[] args) throws Exception { ... SourceFunction<String> source = createSource(applicationParameters.get("InputStream0")); DataStream<String> input = env.addSource(source, "Kinesis Source"); ... }
    • Die Anwendung definiert dann mithilfe des Kinesis Streams Sink-Connectors eine Senke, um Daten an den Ausgabestream zu senden. Der Name des Ausgabestreams und die Region sind in PropertyGroupId = definiertOutputStream0, ähnlich wie beim Eingabestream. Die Senke ist direkt mit der internen Senke verbundenDataStream, die Daten von der Quelle erhält. In einer echten Anwendung findet eine gewisse Transformation zwischen Quelle und Senke statt.

      private static KinesisStreamsSink<String> createSink(Properties outputProperties) { String outputStreamName = outputProperties.getProperty("stream.name"); return KinesisStreamsSink.<String>builder() .setKinesisClientProperties(outputProperties) .setSerializationSchema(new SimpleStringSchema()) .setStreamName(outputStreamName) .setPartitionKeyGenerator(element -> String.valueOf(element.hashCode())) .build(); } ... public static void main(String[] args) throws Exception { ... Sink<String> sink = createSink(applicationParameters.get("OutputStream0")); input.sinkTo(sink); ... }
    • Schließlich führen Sie den Datenfluss aus, den Sie gerade definiert haben. Dies muss die letzte Anweisung der main() Methode sein, nachdem Sie alle Operatoren definiert haben, die der Datenfluss benötigt:

      env.execute("Flink streaming Java API skeleton");

Verwenden Sie die Datei pom.xml

Die Datei pom.xml definiert alle von der Anwendung benötigten Abhängigkeiten und richtet das Maven Shade-Plugin ein, um das Fat-Jar zu erstellen, das alle von Flink benötigten Abhängigkeiten enthält.

  • Einige Abhängigkeiten haben einen Geltungsbereich. provided Diese Abhängigkeiten sind automatisch verfügbar, wenn die Anwendung in Amazon Managed Service für Apache Flink ausgeführt wird. Sie sind erforderlich, um die Anwendung zu kompilieren oder um die Anwendung lokal in Ihrer IDE auszuführen. Weitere Informationen finden Sie unter Führen Sie Ihre Anwendung lokal aus. Stellen Sie sicher, dass Sie dieselbe Flink-Version wie die Runtime verwenden, die Sie in Amazon Managed Service für Apache Flink verwenden werden.

    <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-clients</artifactId> <version>${flink.version}</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>${flink.version}</version> <scope>provided</scope> </dependency>
  • Sie müssen dem POM zusätzliche Apache Flink-Abhängigkeiten mit dem Standardbereich hinzufügen, z. B. den von dieser Anwendung verwendeten Kinesis-Connector. Weitere Informationen finden Sie unter Verwenden Sie Apache Flink-Konnektoren. Sie können auch alle zusätzlichen Java-Abhängigkeiten hinzufügen, die für Ihre Anwendung erforderlich sind.

    <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kinesis</artifactId> <version>${aws.connector.version}</version> </dependency>
  • Das Maven Java Compiler-Plugin stellt sicher, dass der Code mit Java 11 kompiliert wird, der JDK-Version, die derzeit von Apache Flink unterstützt wird.

  • Das Maven Shade-Plugin packt das Fat-Jar, mit Ausnahme einiger Bibliotheken, die von der Runtime bereitgestellt werden. Es spezifiziert auch zwei Transformatoren: und. ServicesResourceTransformer ManifestResourceTransformer Letzteres konfiguriert die Klasse, die die main Methode zum Starten der Anwendung enthält. Wenn Sie die Hauptklasse umbenennen, vergessen Sie nicht, diesen Transformator zu aktualisieren.

  • <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> ... <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer"> <mainClass>com.amazonaws.services.msf.BasicStreamingJob</mainClass> </transformer> ... </plugin>

Schreiben Sie Beispieldatensätze in den Eingabestream

In diesem Abschnitt senden Sie Beispieldatensätze an den Stream, damit die Anwendung sie verarbeiten kann. Sie haben zwei Möglichkeiten, Beispieldaten zu generieren, entweder mit einem Python-Skript oder mit dem Kinesis Data Generator.

Generieren Sie Beispieldaten mithilfe eines Python-Skripts

Sie können ein Python-Skript verwenden, um Beispieldatensätze an den Stream zu senden.

Anmerkung

Um dieses Python-Skript auszuführen, müssen Sie Python 3.x verwenden und die AWS SDK for Python (Boto) -Bibliothek installiert haben.

Um mit dem Senden von Testdaten an den Kinesis-Eingabestream zu beginnen:

  1. Laden Sie das stock.py Python-Skript für den Datengenerator aus dem GitHub Datengenerator-Repository herunter.

  2. Führen Sie das stock.pySkript aus:

    $ python stock.py

Lassen Sie das Skript laufen, während Sie den Rest des Tutorials abschließen. Sie können jetzt Ihre Apache Flink-Anwendung ausführen.

Generieren Sie Beispieldaten mit dem Kinesis Data Generator

Alternativ zur Verwendung des Python-Skripts können Sie den Kinesis Data Generator verwenden, der auch in einer gehosteten Version verfügbar ist, um Zufallsstichprobendaten an den Stream zu senden. Kinesis Data Generator wird in Ihrem Browser ausgeführt, und Sie müssen nichts auf Ihrem Computer installieren.

So richten Sie Kinesis Data Generator ein und führen ihn aus:

  1. Folgen Sie den Anweisungen in der Kinesis Data Generator-Dokumentation, um den Zugriff auf das Tool einzurichten. Sie führen eine CloudFormation Vorlage aus, die einen Benutzer und ein Passwort einrichtet.

  2. Greifen Sie über die von der CloudFormation Vorlage generierte URL auf Kinesis Data Generator zu. Sie finden die URL auf der Registerkarte „Ausgabe“, nachdem die CloudFormation Vorlage fertiggestellt ist.

  3. Konfigurieren Sie den Datengenerator:

    • Region: Wählen Sie die Region aus, die Sie für dieses Tutorial verwenden: us-east-1

    • Stream/delivery Stream: Wählen Sie den Eingabestream aus, den die Anwendung verwenden soll: ExampleInputStream

    • Aufzeichnungen pro Sekunde: 100

    • Datensatzvorlage: Kopieren Sie die folgende Vorlage und fügen Sie sie ein:

      { "event_time" : "{{date.now("YYYY-MM-DDTkk:mm:ss.SSSSS")}}, "ticker" : "{{random.arrayElement( ["AAPL", "AMZN", "MSFT", "INTC", "TBV"] )}}", "price" : {{random.number(100)}} }
  4. Testen Sie die Vorlage: Wählen Sie Vorlage testen und überprüfen Sie, ob der generierte Datensatz dem folgenden ähnelt:

    { "event_time" : "2024-06-12T15:08:32.04800, "ticker" : "INTC", "price" : 7 }
  5. Starten Sie den Datengenerator: Wählen Sie „Daten senden“.

Kinesis Data Generator sendet jetzt Daten an dieExampleInputStream.

Führen Sie Ihre Anwendung lokal aus

Sie können Ihre Flink-Anwendung lokal in Ihrer IDE ausführen und debuggen.

Anmerkung

Bevor Sie fortfahren, stellen Sie sicher, dass die Eingabe- und Ausgabestreams verfügbar sind. Siehe Erstellen Sie zwei Amazon Kinesis-Datenstreams. Stellen Sie außerdem sicher, dass Sie über Lese- und Schreibberechtigungen für beide Streams verfügen. Siehe Authentifizieren Sie Ihre AWS Sitzung.

Für die Einrichtung der lokalen Entwicklungsumgebung sind Java 11 JDK, Apache Maven und eine IDE für die Java-Entwicklung erforderlich. Stellen Sie sicher, dass Sie die erforderlichen Voraussetzungen erfüllen. Siehe Erfüllen Sie die Voraussetzungen für den Abschluss der Übungen.

Importieren Sie das Java-Projekt in Ihre IDE

Um mit der Arbeit an der Anwendung in Ihrer IDE zu beginnen, müssen Sie sie als Java-Projekt importieren.

Das Repository, das Sie geklont haben, enthält mehrere Beispiele. Jedes Beispiel ist ein separates Projekt. Importieren Sie für dieses Tutorial den Inhalt des ./java/GettingStarted Unterverzeichnisses in Ihre IDE.

Fügen Sie den Code als vorhandenes Java-Projekt mit Maven ein.

Anmerkung

Der genaue Vorgang zum Importieren eines neuen Java-Projekts hängt von der verwendeten IDE ab.

Überprüfen Sie die lokale Anwendungskonfiguration

Wenn die Anwendung lokal ausgeführt wird, verwendet sie die Konfiguration in der application_properties.json Datei im Ressourcenordner des Projekts unter./src/main/resources. Sie können diese Datei bearbeiten, um andere Kinesis-Stream-Namen oder -Regionen zu verwenden.

[ { "PropertyGroupId": "InputStream0", "PropertyMap": { "stream.name": "ExampleInputStream", "flink.stream.initpos": "LATEST", "aws.region": "us-east-1" } }, { "PropertyGroupId": "OutputStream0", "PropertyMap": { "stream.name": "ExampleOutputStream", "aws.region": "us-east-1" } } ]

Richten Sie Ihre IDE-Ausführungskonfiguration ein

Sie können die Flink-Anwendung direkt von Ihrer IDE aus ausführen und debuggen, indem Sie die Hauptklasse ausführencom.amazonaws.services.msf.BasicStreamingJob, wie Sie jede Java-Anwendung ausführen würden. Bevor Sie die Anwendung ausführen, müssen Sie die Run-Konfiguration einrichten. Das Setup hängt von der IDE ab, die Sie verwenden. Sehen Sie sich beispielsweise Run/debug Konfigurationen in der IntelliJ IDEA-Dokumentation an. Insbesondere müssen Sie Folgendes einrichten:

  1. Fügen Sie die provided Abhängigkeiten zum Klassenpfad hinzu. Dies ist erforderlich, um sicherzustellen, dass die Abhängigkeiten mit provided Gültigkeitsbereich an die Anwendung weitergegeben werden, wenn sie lokal ausgeführt werden. Ohne diese Einrichtung zeigt die Anwendung sofort einen class not found Fehler an.

  2. Übergeben Sie die AWS Anmeldeinformationen für den Zugriff auf die Kinesis-Streams an die Anwendung. Der schnellste Weg ist die Verwendung von AWS Toolkit für IntelliJ IDEA. Mit diesem IDE-Plugin in der Run-Konfiguration können Sie ein bestimmtes Profil auswählen. AWS AWS Die Authentifizierung erfolgt mit diesem Profil. Sie müssen die AWS Anmeldeinformationen nicht direkt weitergeben.

  3. Stellen Sie sicher, dass die IDE die Anwendung mit JDK 11 ausführt.

Führen Sie die Anwendung in Ihrer IDE aus

Nachdem Sie die Run-Konfiguration für den eingerichtet habenBasicStreamingJob, können Sie sie wie eine normale Java-Anwendung ausführen oder debuggen.

Anmerkung

Sie können das von Maven generierte Fat-Jar nicht direkt über die java -jar ... Befehlszeile ausführen. Dieses Jar enthält nicht die Flink-Kernabhängigkeiten, die für die eigenständige Ausführung der Anwendung erforderlich sind.

Wenn die Anwendung erfolgreich gestartet wird, protokolliert sie einige Informationen über den eigenständigen Minicluster und die Initialisierung der Konnektoren. Darauf folgen eine Reihe von INFO- und einige WARN-Logs, die Flink normalerweise ausgibt, wenn die Anwendung gestartet wird.

13:43:31,405 INFO com.amazonaws.services.msf.BasicStreamingJob [] - Loading application properties from 'flink-application-properties-dev.json' 13:43:31,549 INFO org.apache.flink.streaming.connectors.kinesis.FlinkKinesisConsumer [] - Flink Kinesis Consumer is going to read the following streams: ExampleInputStream, 13:43:31,676 INFO org.apache.flink.runtime.taskexecutor.TaskExecutorResourceUtils [] - The configuration option taskmanager.cpu.cores required for local execution is not set, setting it to the maximal possible value. 13:43:31,676 INFO org.apache.flink.runtime.taskexecutor.TaskExecutorResourceUtils [] - The configuration option taskmanager.memory.task.heap.size required for local execution is not set, setting it to the maximal possible value. 13:43:31,676 INFO org.apache.flink.runtime.taskexecutor.TaskExecutorResourceUtils [] - The configuration option taskmanager.memory.task.off-heap.size required for local execution is not set, setting it to the maximal possible value. 13:43:31,676 INFO org.apache.flink.runtime.taskexecutor.TaskExecutorResourceUtils [] - The configuration option taskmanager.memory.network.min required for local execution is not set, setting it to its default value 64 mb. 13:43:31,676 INFO org.apache.flink.runtime.taskexecutor.TaskExecutorResourceUtils [] - The configuration option taskmanager.memory.network.max required for local execution is not set, setting it to its default value 64 mb. 13:43:31,676 INFO org.apache.flink.runtime.taskexecutor.TaskExecutorResourceUtils [] - The configuration option taskmanager.memory.managed.size required for local execution is not set, setting it to its default value 128 mb. 13:43:31,677 INFO org.apache.flink.runtime.minicluster.MiniCluster [] - Starting Flink Mini Cluster ....

Nach Abschluss der Initialisierung gibt die Anwendung keine weiteren Logeinträge aus. Während des Datenflusses wird kein Protokoll ausgegeben.

Um zu überprüfen, ob die Anwendung Daten korrekt verarbeitet, können Sie die Eingabe- und Ausgabe-Kinesis-Streams überprüfen, wie im folgenden Abschnitt beschrieben.

Anmerkung

Das normale Verhalten einer Flink-Anwendung besteht darin, keine Protokolle über den Datenfluss auszugeben. Das Ausgeben von Protokollen für jeden Datensatz mag für das Debuggen praktisch sein, kann aber bei der Ausführung in der Produktion zu erheblichem Aufwand führen.

Beobachten Sie Eingabe- und Ausgabedaten in Kinesis-Streams

Mithilfe des Data Viewers in der Amazon Kinesis-Konsole können Sie Datensätze beobachten, die vom (generierenden Beispiel-Python) oder vom Kinesis Data Generator (Link) an den Eingabestream gesendet wurden.

Um Aufzeichnungen zu beobachten
  1. Öffnen Sie die Kinesis-Konsole unter https://console.aws.amazon.com/kinesis.

  2. Stellen Sie sicher, dass die Region mit der Region übereinstimmt, in der Sie dieses Tutorial ausführen. Die Standardeinstellung lautet us-east-1 US East (Nord-Virginia). Ändern Sie die Region, wenn sie nicht übereinstimmt.

  3. Wählen Sie Datenströme.

  4. Wählen Sie den Stream aus, den Sie beobachten möchten, entweder ExampleInputStream oder ExampleOutputStream.

  5. Wählen Sie den Tab „Datenanzeige“.

  6. Wählen Sie einen beliebigen Shard aus, behalten Sie „Neueste“ als Startposition bei und wählen Sie dann „Datensätze abrufen“. Möglicherweise wird die Fehlermeldung „Für diese Anfrage wurde kein Datensatz gefunden“ angezeigt. Wenn ja, wählen Sie Erneut versuchen, Datensätze abzurufen. Die neuesten im Stream veröffentlichten Datensätze werden angezeigt.

  7. Wählen Sie den Wert in der Datenspalte aus, um den Inhalt des Datensatzes im JSON-Format zu überprüfen.

Beenden Sie die lokale Ausführung Ihrer Anwendung

Stoppen Sie die Ausführung der Anwendung in Ihrer IDE. Die IDE bietet normalerweise eine „Stopp“ -Option. Der genaue Ort und die Methode hängen von der verwendeten IDE ab.

Kompilieren und verpacken Sie Ihren Anwendungscode

In diesem Abschnitt verwenden Sie Apache Maven, um den Java-Code zu kompilieren und in eine JAR-Datei zu packen. Sie können Ihren Code mit dem Maven-Befehlszeilentool oder Ihrer IDE kompilieren und paketieren.

Um mit der Maven-Befehlszeile zu kompilieren und zu paketieren:

Gehen Sie in das Verzeichnis, das das GettingStarted Java-Projekt enthält, und führen Sie den folgenden Befehl aus:

$ mvn package

Um mit Ihrer IDE zu kompilieren und zu paketieren:

Führen Sie es mvn package von Ihrer IDE Maven-Integration aus.

In beiden Fällen wird die folgende JAR-Datei erstellt:target/amazon-msf-java-stream-app-1.0.jar.

Anmerkung

Wenn Sie ein „Build-Projekt“ von Ihrer IDE aus ausführen, wird die JAR-Datei möglicherweise nicht erstellt.

Laden Sie die JAR-Datei mit dem Anwendungscode hoch

In diesem Abschnitt laden Sie die JAR-Datei, die Sie im vorherigen Abschnitt erstellt haben, in den Amazon Simple Storage Service (Amazon S3) -Bucket hoch, den Sie zu Beginn dieses Tutorials erstellt haben. Wenn Sie diesen Schritt noch nicht abgeschlossen haben, finden Sie weitere Informationen unter (Link).

Um die JAR-Datei mit dem Anwendungscode hochzuladen
  1. Öffnen Sie die Amazon S3-Konsole unter https://console.aws.amazon.com/s3/.

  2. Wählen Sie den Bucket aus, den Sie zuvor für den Anwendungscode erstellt haben.

  3. Klicken Sie auf Upload.

  4. Klicken Sie auf Add files.

  5. Navigieren Sie zu der JAR-Datei, die im vorherigen Schritt generiert wurde:target/amazon-msf-java-stream-app-1.0.jar.

  6. Wählen Sie Hochladen, ohne weitere Einstellungen zu ändern.

Warnung

Stellen Sie sicher, dass Sie die richtige JAR-Datei in auswählen<repo-dir>/java/GettingStarted/target/amazon-msf-java-stream-app-1.0.jar.

Das target Verzeichnis enthält auch andere JAR-Dateien, die Sie nicht hochladen müssen.

Erstellen und konfigurieren Sie die Anwendung Managed Service für Apache Flink

Sie können eine Anwendung von Managed Service für Apache Flink entweder über die Konsole oder AWS CLI erstellen und ausführen. Für dieses Tutorial verwenden Sie die Konsole.

Anmerkung

Wenn Sie die Anwendung mithilfe der Konsole erstellen, werden Ihre AWS Identity and Access Management (IAM) und Amazon CloudWatch Logs-Ressourcen für Sie erstellt. Wenn Sie die Anwendung mit der erstellen AWS CLI, erstellen Sie diese Ressourcen separat.

Erstellen der Anwendung

So erstellen Sie die Anwendung
  1. Melden Sie sich bei der AWS-Managementkonsole an und öffnen Sie die Amazon MSF-Konsole unter https://console.aws.amazon.com/flink.

  2. Stellen Sie sicher, dass die richtige Region ausgewählt ist: us-east-1 US East (Nord-Virginia)

  3. Öffnen Sie das Menü auf der rechten Seite und wählen Sie Apache Flink-Anwendungen und dann Streaming-Anwendung erstellen. Wählen Sie alternativ im Container Erste Schritte auf der Startseite die Option Streaming-Anwendung erstellen.

  4. Gehen Sie auf der Seite „Streaming-Anwendung erstellen“ wie folgt vor:

    • Wählen Sie eine Methode zum Einrichten der Stream-Verarbeitungsanwendung: Wählen Sie „Von Grund auf neu erstellen“.

    • Apache Flink-Konfiguration, Flink-Version der Anwendung: Wählen Sie Apache Flink 1.20.

  5. Konfigurieren Sie Ihre Anwendung

    • Name der Anwendung: geben Sie einMyApplication.

    • Beschreibung: eingebenMy java test app.

    • Zugriff auf Anwendungsressourcen: Wählen Sie IAM-Rolle kinesis-analytics-MyApplication-us-east-1 mit den erforderlichen Richtlinien erstellen/aktualisieren.

  6. Konfigurieren Sie Ihre Vorlage für Anwendungseinstellungen

    • Vorlagen: Wählen Sie Entwicklung.

  7. Wählen Sie unten auf der Seite Streaming-Anwendung erstellen aus.

Anmerkung

Beim Erstellen einer Anwendung von Managed Service für Apache Flink mit der Konsole haben Sie die Möglichkeit, eine IAM-Rolle und -Richtlinie für Ihre Anwendung erstellen zu lassen. Ihre Anwendung verwendet diese Rolle und Richtlinie für den Zugriff auf ihre abhängigen Ressourcen. Diese IAM-Ressourcen werden unter Verwendung Ihres Anwendungsnamens und der Region wie folgt benannt:

  • Richtlinie: kinesis-analytics-service-MyApplication-us-east-1

  • Rolle: kinesisanalytics-MyApplication-us-east-1

Amazon Managed Service für Apache Flink war früher als Kinesis Data Analytics bekannt. Aus Gründen der Abwärtskompatibilität wird dem Namen der automatisch erstellten Ressourcen ein Präfix kinesis-analytics- vorangestellt.

Bearbeiten Sie die IAM-Richtlinie

Bearbeiten Sie die IAM-Richtlinie zum Hinzufügen von Berechtigungen für den Zugriff auf die Kinesis-Datenströme.

Um die Richtlinie zu bearbeiten
  1. Öffnen Sie unter https://console.aws.amazon.com/iam/ die IAM-Konsole.

  2. Wählen Sie Policies (Richtlinien). Wählen Sie die kinesis-analytics-service-MyApplication-us-east-1-Richtlinie aus, die die Konsole im vorherigen Abschnitt für Sie erstellt hat.

  3. Wählen Sie Bearbeiten und dann die Registerkarte JSON.

  4. Fügen Sie den markierten Abschnitt der folgenden Beispielrichtlinie der Richtlinie hinzu. Ersetzen Sie die Beispielkonto-IDs (012345678901) durch Ihre Konto-ID.

    JSON
    { "Version":"2012-10-17", "Statement": [ { "Sid": "ReadCode", "Effect": "Allow", "Action": [ "s3:GetObject", "s3:GetObjectVersion" ], "Resource": [ "arn:aws:s3:::my-bucket/kinesis-analytics-placeholder-s3-object" ] }, { "Sid": "ListCloudwatchLogGroups", "Effect": "Allow", "Action": [ "logs:DescribeLogGroups" ], "Resource": [ "arn:aws:logs:us-east-1:012345678901:log-group:*" ] }, { "Sid": "ListCloudwatchLogStreams", "Effect": "Allow", "Action": [ "logs:DescribeLogStreams" ], "Resource": [ "arn:aws:logs:us-east-1:012345678901:log-group:/aws/kinesis-analytics/MyApplication:log-stream:*" ] }, { "Sid": "PutCloudwatchLogs", "Effect": "Allow", "Action": [ "logs:PutLogEvents" ], "Resource": [ "arn:aws:logs:us-east-1:012345678901:log-group:/aws/kinesis-analytics/MyApplication:log-stream:kinesis-analytics-log-stream" ] }, { "Sid": "ReadInputStream", "Effect": "Allow", "Action": "kinesis:*", "Resource": "arn:aws:kinesis:us-east-1:012345678901:stream/ExampleInputStream" }, { "Sid": "WriteOutputStream", "Effect": "Allow", "Action": "kinesis:*", "Resource": "arn:aws:kinesis:us-east-1:012345678901:stream/ExampleOutputStream" } ] }
  5. Wählen Sie unten auf der Seite Weiter und dann Änderungen speichern aus.

Konfigurieren Sie die Anwendung

Bearbeiten Sie die Anwendungskonfiguration, um das Anwendungscode-Artefakt festzulegen.

Um die Konfiguration zu bearbeiten
  1. Wählen Sie auf der MyApplication Seite „Konfigurieren“.

  2. Gehen Sie im Abschnitt Speicherort des Anwendungscodes wie folgt vor:

    • Wählen Sie für den Amazon S3-Bucket den Bucket aus, den Sie zuvor für den Anwendungscode erstellt haben. Wählen Sie Browse und wählen Sie den richtigen Bucket aus. Wählen Sie dann Choose aus. Klicken Sie nicht auf den Bucket-Namen.

    • Geben Sie als Pfad zum Amazon-S3-Objekt den Wert amazon-msf-java-stream-app-1.0.jar ein.

  3. Wählen Sie für Zugriffsberechtigungen die Option IAM-Rolle kinesis-analytics-MyApplication-us-east-1 mit den erforderlichen Richtlinien erstellen/aktualisieren aus.

  4. Fügen Sie im Abschnitt Runtime-Eigenschaften die folgenden Eigenschaften hinzu.

  5. Wählen Sie Neues Objekt hinzufügen und fügen Sie jeden der folgenden Parameter hinzu:

    Gruppen-ID Key (Schlüssel) Value (Wert)
    InputStream0 stream.name ExampleInputStream
    InputStream0 aws.region us-east-1
    OutputStream0 stream.name ExampleOutputStream
    OutputStream0 aws.region us-east-1
  6. Ändern Sie keinen der anderen Abschnitte.

  7. Wählen Sie Änderungen speichern aus.

Anmerkung

Wenn Sie Amazon CloudWatch Logging aktivieren, erstellt Managed Service for Apache Flink eine Log-Gruppe und einen Log-Stream für Sie. Die Namen dieser Ressourcen lauten wie folgt:

  • Protokollgruppe: /aws/kinesis-analytics/MyApplication

  • Protokollstream: kinesis-analytics-log-stream

Führen Sie die Anwendung aus.

Die Anwendung ist jetzt konfiguriert und kann ausgeführt werden.

Ausführen der Anwendung
  1. Wählen Sie auf der Konsole für Amazon Managed Service für Apache Flink Meine Anwendung und dann Ausführen aus.

  2. Wählen Sie auf der nächsten Seite, der Konfigurationsseite zur Anwendungswiederherstellung, die Option Mit dem neuesten Snapshot ausführen und dann Ausführen aus.

    Der Status in den Anwendungsdetails wechselt von Ready zu Starting und dann zu dem Running Zeitpunkt, zu dem die Anwendung gestartet wurde.

Wenn sich die Anwendung im Running Status befindet, können Sie jetzt das Flink-Dashboard öffnen.

So öffnen Sie das -Dashboard
  1. Wählen Sie Apache Flink-Dashboard öffnen. Das Dashboard wird auf einer neuen Seite geöffnet.

  2. Wählen Sie in der Liste „Laufende Jobs“ den einzelnen Job aus, den Sie sehen können.

    Anmerkung

    Wenn Sie die Runtime-Eigenschaften festgelegt oder die IAM-Richtlinien falsch bearbeitet haben, wird der Anwendungsstatus möglicherweise wie folgt geändertRunning, aber das Flink-Dashboard zeigt an, dass der Job kontinuierlich neu gestartet wird. Dies ist ein häufiges Fehlerszenario, wenn die Anwendung falsch konfiguriert ist oder nicht über Berechtigungen für den Zugriff auf die externen Ressourcen verfügt.

    In diesem Fall überprüfen Sie die Registerkarte Ausnahmen im Flink-Dashboard, um die Ursache des Problems zu ermitteln.

Beachten Sie die Metriken der laufenden Anwendung

Auf der MyApplication Seite finden Sie im Abschnitt CloudWatch Amazon-Metriken einige der grundlegenden Metriken der laufenden Anwendung.

Um die Metriken einzusehen
  1. Wählen Sie neben der Schaltfläche „Aktualisieren“ 10 Sekunden aus der Dropdownliste aus.

  2. Wenn die Anwendung läuft und fehlerfrei ist, können Sie sehen, dass die Verfügbarkeitsmetrik kontinuierlich ansteigt.

  3. Die Metrik für vollständige Neustarts sollte Null sein. Wenn sie zunimmt, kann es bei der Konfiguration zu Problemen kommen. Um das Problem zu untersuchen, überprüfen Sie die Registerkarte Ausnahmen im Flink-Dashboard.

  4. Die Metrik „Anzahl der fehlgeschlagenen Checkpoints“ sollte in einer fehlerfreien Anwendung Null sein.

    Anmerkung

    Dieses Dashboard zeigt einen festen Satz von Metriken mit einer Granularität von 5 Minuten an. Sie können ein benutzerdefiniertes Anwendungs-Dashboard mit beliebigen Metriken im CloudWatch Dashboard erstellen.

Beobachten Sie die Ausgabedaten in Kinesis-Streams

Stellen Sie sicher, dass Sie weiterhin Daten für die Eingabe veröffentlichen, entweder mithilfe des Python-Skripts oder des Kinesis-Datengenerators.

Sie können jetzt die Ausgabe der Anwendung beobachten, die auf Managed Service for Apache Flink ausgeführt wird, indem Sie den Data Viewer in der verwenden https://console.aws.amazon.com/kinesis/, ähnlich wie Sie es bereits zuvor getan haben.

Um die Ausgabe anzusehen
  1. Öffnen Sie die Kinesis-Konsole unter https://console.aws.amazon.com/kinesis.

  2. Stellen Sie sicher, dass die Region mit der Region übereinstimmt, die Sie zum Ausführen dieses Tutorials verwenden. In der Standardeinstellung ist es US-East-1US East (Nord-Virginia). Ändern Sie bei Bedarf die Region.

  3. Wählen Sie Datenströme.

  4. Wählen Sie den Stream aus, den Sie beobachten möchten. Verwenden Sie für dieses Tutorial ExampleOutputStream.

  5. Wählen Sie die Registerkarte Datenanzeige.

  6. Wählen Sie einen beliebigen Shard aus, behalten Sie „Neueste“ als Startposition bei und wählen Sie dann „Datensätze abrufen“. Möglicherweise wird die Fehlermeldung „Für diese Anfrage wurde kein Datensatz gefunden“ angezeigt. Wenn ja, wählen Sie Erneut versuchen, Datensätze abzurufen. Die neuesten im Stream veröffentlichten Datensätze werden angezeigt.

  7. Wählen Sie den Wert in der Datenspalte aus, um den Inhalt des Datensatzes im JSON-Format zu überprüfen.

Stoppen Sie die Anwendung

Um die Anwendung zu beenden, rufen Sie die Konsolenseite der Anwendung Managed Service for Apache Flink mit dem Namen auf. MyApplication

So stoppen Sie die Anwendung
  1. Wählen Sie in der Dropdownliste Aktion die Option Stopp aus.

  2. Der Status in den Anwendungsdetails wechselt von Running zu und dann zu dem Ready ZeitpunktStopping, an dem die Anwendung vollständig gestoppt wurde.

    Anmerkung

    Vergessen Sie nicht, auch das Senden von Daten aus dem Python-Skript oder dem Kinesis-Datengenerator an den Eingabestream zu beenden.

Nächster Schritt

Bereinigen AWS Ressourcen