This is an automated email from the ASF dual-hosted git repository.
Amar3tto pushed a commit to branch snowflakeio-yaml
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/snowflakeio-yaml by this push:
new c2d8c71af61 Fix bytes
c2d8c71af61 is described below
commit c2d8c71af61eb81f77c2822d1c795cb2c4c438de
Author: Vitaly Terentyev <[email protected]>
AuthorDate: Thu Aug 13 14:47:46 2026 +0400
Fix bytes
---
sdks/java/io/snowflake/build.gradle | 1 +
.../snowflake/SnowflakeSchemaTransformUtils.java | 24 ++++++++++++++++++++--
.../SnowflakeReadSchemaTransformProviderTest.java | 2 +-
3 files changed, 24 insertions(+), 3 deletions(-)
diff --git a/sdks/java/io/snowflake/build.gradle
b/sdks/java/io/snowflake/build.gradle
index 951af328fed..c1a8f2b93ea 100644
--- a/sdks/java/io/snowflake/build.gradle
+++ b/sdks/java/io/snowflake/build.gradle
@@ -31,6 +31,7 @@ dependencies {
permitUnusedDeclared project(path:
":sdks:java:extensions:google-cloud-platform-core")
implementation library.java.slf4j_api
implementation library.java.everit_json_schema
+ permitUnusedDeclared library.java.everit_json_schema
implementation group: 'net.snowflake', name: 'snowflake-jdbc', version:
'4.0.2'
implementation group: 'com.opencsv', name: 'opencsv', version: '5.12.0'
implementation 'net.snowflake:snowflake-ingest-sdk:4.4.2'
diff --git
a/sdks/java/io/snowflake/src/main/java/org/apache/beam/sdk/io/snowflake/SnowflakeSchemaTransformUtils.java
b/sdks/java/io/snowflake/src/main/java/org/apache/beam/sdk/io/snowflake/SnowflakeSchemaTransformUtils.java
index 84d451552c3..7f46071e44c 100644
---
a/sdks/java/io/snowflake/src/main/java/org/apache/beam/sdk/io/snowflake/SnowflakeSchemaTransformUtils.java
+++
b/sdks/java/io/snowflake/src/main/java/org/apache/beam/sdk/io/snowflake/SnowflakeSchemaTransformUtils.java
@@ -17,7 +17,6 @@
*/
package org.apache.beam.sdk.io.snowflake;
-import java.nio.charset.StandardCharsets;
import javax.annotation.Nullable;
import org.apache.beam.sdk.io.snowflake.data.SnowflakeColumn;
import org.apache.beam.sdk.io.snowflake.data.SnowflakeDataType;
@@ -246,7 +245,7 @@ public class SnowflakeSchemaTransformUtils {
throw new IllegalArgumentException(String.format("Invalid boolean
value '%s'.", value));
case BYTES:
- return value.getBytes(StandardCharsets.UTF_8);
+ return decodeHex(value);
case DATETIME:
return Instant.parse(value);
@@ -269,6 +268,27 @@ public class SnowflakeSchemaTransformUtils {
}
}
+ private static byte[] decodeHex(String value) {
+ if ((value.length() & 1) != 0) {
+ throw new IllegalArgumentException("Invalid hexadecimal Snowflake binary
value.");
+ }
+
+ byte[] result = new byte[value.length() / 2];
+
+ for (int i = 0; i < value.length(); i += 2) {
+ int high = Character.digit(value.charAt(i), 16);
+ int low = Character.digit(value.charAt(i + 1), 16);
+
+ if (high == -1 || low == -1) {
+ throw new IllegalArgumentException("Invalid hexadecimal Snowflake
binary value.");
+ }
+
+ result[i / 2] = (byte) ((high << 4) | low);
+ }
+
+ return result;
+ }
+
public static SnowflakeDataType toSnowflakeDataType(Schema.Field field) {
switch (field.getType().getTypeName()) {
case BYTE:
diff --git
a/sdks/java/io/snowflake/src/test/java/org/apache/beam/sdk/io/snowflake/SnowflakeReadSchemaTransformProviderTest.java
b/sdks/java/io/snowflake/src/test/java/org/apache/beam/sdk/io/snowflake/SnowflakeReadSchemaTransformProviderTest.java
index 02f8e2e7b66..440103b7cfb 100644
---
a/sdks/java/io/snowflake/src/test/java/org/apache/beam/sdk/io/snowflake/SnowflakeReadSchemaTransformProviderTest.java
+++
b/sdks/java/io/snowflake/src/test/java/org/apache/beam/sdk/io/snowflake/SnowflakeReadSchemaTransformProviderTest.java
@@ -124,7 +124,7 @@ public class SnowflakeReadSchemaTransformProviderTest {
.build();
String[] values = {
- "1", "2", "3", "4", "5.5", "6.5", "hello", "true", "abc",
"2026-08-13T09:00:00.000Z"
+ "1", "2", "3", "4", "5.5", "6.5", "hello", "true", "616263",
"2026-08-13T09:00:00.000Z"
};
Row row = SnowflakeSchemaTransformUtils.toRow(values, schema);