Skip to content

Byzer-Python

This library is an ongoing effort towards bringing the data exchanging ability between Java/Scala and Python. PyJava introduces Apache Arrow as the exchanging data format, this means we can avoid ser/der between Java/Scala and Python which can really speed up the communication efficiency than traditional way.

When you invoke python code in Java/Scala side, PyJava will start some python workers automatically and send the data to python worker, and once they are processed, send them back. The python workers are reused
by default.

The initial code in this lib is from Apache Spark.

Install

Setup python(>= 3.6) Env(Conda is recommended):

pip uninstall pyjava && pip install pyjava

Setup Java env(Maven is recommended):

For Scala 2.11/Spark 2.4.3

<dependency>
<groupId>tech.mlsql</groupId>
<artifactId>pyjava-2.4_2.11</artifactId>
<version>0.3.0</version>
</dependency>

For Scala 2.12/Spark 3.0.1

<dependency>
<groupId>tech.mlsql</groupId>
<artifactId>pyjava-3.0_2.12</artifactId>
<version>0.3.0</version>
</dependency>

For Scala 2.12/Spark 3.1.1

<dependency>
<groupId>tech.mlsql</groupId>
<artifactId>pyjava-3.0_2.12</artifactId>
<version>0.3.2</version>
</dependency>

Build Mannually

Install Build Tool:

pip install mlsql_plugin_tool

Build for Spark 3.1.1:

mlsql_plugin_tool spark311
mvn clean install -DskipTests -Pdisable-java8-doclint -Prelease-sign-artifacts

Build For Spark 2.4.3

mlsql_plugin_tool spark243
mvn clean install -DskipTests -Pdisable-java8-doclint -Prelease-sign-artifacts

Using python code snippet to process data in Java/Scala

With pyjava, you can run any python code in your Java/Scala application.

valenvs=new util.HashMap[String, String]()
// prepare python environment
envs.put(str(PythonConf.PYTHON_ENV), "source activate dev && export ARROW_PRE_0_15_IPC_FORMAT=1 ")
// describe the data which will be transfered to python valsourceSchema=StructType(Seq(StructField("value", StringType)))
valbatch=newArrowPythonRunner(
Seq(ChainedPythonFunctions(Seq(PythonFunction(
""" |import pandas as pd |import numpy as np | |def process(): | for item in context.fetch_once_as_rows(): | item["value1"] = item["value"] + "_suffix" | yield item | |context.build_result(process())""".stripMargin, envs, "python", "3.6")))), sourceSchema,
"GMT", Map()
)
// prepare datavalsourceEnconder=RowEncoder.apply(sourceSchema).resolveAndBind()
valnewIter=Seq(Row.fromSeq(Seq("a1")), Row.fromSeq(Seq("a2"))).map { irow =>
sourceEnconder.toRow(irow).copy()
}.iterator
// run the code and get the return resultvaljavaConext=newJavaContextvalcommonTaskContext=newAppContextImpl(javaConext, batch)
valcolumnarBatchIter= batch.compute(Iterator(newIter), TaskContext.getPartitionId(), commonTaskContext)
//f.copy(), copy function is required 
columnarBatchIter.flatMap { batch =>
batch.rowIterator.asScala
}.foreach(f => println(f.copy()))
javaConext.markComplete
javaConext.close

Using python code snippet to process data in Spark

valsession= spark
importsession.implicits._valtimezoneid= session.sessionState.conf.sessionLocalTimeZone
valdf= session.createDataset[String](Seq("a1", "b1")).toDF("value")
valstruct= df.schema
valabc= df.rdd.mapPartitions { iter =>valenconder=RowEncoder.apply(struct).resolveAndBind()
valenvs=new util.HashMap[String, String]()
envs.put(str(PythonConf.PYTHON_ENV), "source activate streamingpro-spark-2.4.x")
valbatch=newArrowPythonRunner(
Seq(ChainedPythonFunctions(Seq(PythonFunction(
""" |import pandas as pd |import numpy as np |for item in data_manager.fetch_once(): | print(item) |df = pd.DataFrame({'AAA': [4, 5, 6, 7],'BBB': [10, 20, 30, 40],'CCC': [100, 50, -30, -50]}) |data_manager.set_output([[df['AAA'],df['BBB']]])""".stripMargin, envs, "python", "3.6")))), struct,
timezoneid, Map()
)
valnewIter= iter.map { irow =>
enconder.toRow(irow)
}
valcommonTaskContext=newSparkContextImp(TaskContext.get(), batch)
valcolumnarBatchIter= batch.compute(Iterator(newIter), TaskContext.getPartitionId(), commonTaskContext)
columnarBatchIter.flatMap { batch =>
batch.rowIterator.asScala.map(_.copy)
}
}
valwow=SparkUtils.internalCreateDataFrame(session, abc, StructType(Seq(StructField("AAA", LongType), StructField("BBB", LongType))), false)
wow.show()

Run Python Project

With Pyjava, you can tell the system where is the python project and which is then entrypoint, then you can run this project in Java/Scala.

importtech.mlsql.arrow.python.runner.PythonProjectRunnervalrunner=newPythonProjectRunner("./pyjava/examples/pyproject1", Map())
valoutput= runner.run(Seq("bash", "-c", "source activate dev && python train.py"), Map(
"tempDataLocalPath"->"/tmp/data",
"tempModelLocalPath"->"/tmp/model"
))
output.foreach(println)

Example In MLSQL

None Interactive Mode:

!python env "PYTHON_ENV=source activate streamingpro-spark-2.4.x";
!python conf "schema=st(field(a,long),field(b,long))";
select1as a as table1;
!python on table1 '''import pandas as pdimport numpy as npfor item in data_manager.fetch_once(): print(item)df = pd.DataFrame({'AAA': [4, 5, 6, 8],'BBB': [10, 20, 30, 40],'CCC': [100, 50, -30, -50]})data_manager.set_output([[df['AAA'],df['BBB']]])''' named mlsql_temp_table2;
select*from mlsql_temp_table2 as output; 

Interactive Mode:

!python start;
!python env "PYTHON_ENV=source activate streamingpro-spark-2.4.x";
!python env "schema=st(field(a,integer),field(b,integer))";
!python '''import pandas as pdimport numpy as np''';
!python '''for item in data_manager.fetch_once(): print(item)df = pd.DataFrame({'AAA': [4, 5, 6, 8],'BBB': [10, 20, 30, 40],'CCC': [100, 50, -30, -50]})data_manager.set_output([[df['AAA'],df['BBB']]])''';
!python close;

Using PyJava as Arrow Server/Client

Java Server side:

valsocketRunner=newSparkSocketRunner("wow", NetUtils.getHost, "Asia/Harbin")
valdataSchema=StructType(Seq(StructField("value", StringType)))
valenconder=RowEncoder.apply(dataSchema).resolveAndBind()
valnewIter=Seq(Row.fromSeq(Seq("a1")), Row.fromSeq(Seq("a2"))).map { irow =>
enconder.toRow(irow)
}.iterator
valjavaConext=newJavaContextvalcommonTaskContext=newAppContextImpl(javaConext, null)
valArray(_, host, port) = socketRunner.serveToStreamWithArrow(newIter, dataSchema, 10, commonTaskContext)
println(s"${host}:${port}")
Thread.currentThread().join()

Python Client side:

importosimportsocketfrompyjava.serializersimport \
ArrowStreamPandasSerializerout_ser=ArrowStreamPandasSerializer(None, True, True)
out_ser=ArrowStreamPandasSerializer("Asia/Harbin", False, None)
HOST=""PORT=-1withsocket.socket(socket.AF_INET, socket.SOCK_STREAM) assock:
sock.connect((HOST, PORT))
buffer_size=int(os.environ.get("SPARK_BUFFER_SIZE", 65536))
infile=os.fdopen(os.dup(sock.fileno()), "rb", buffer_size)
outfile=os.fdopen(os.dup(sock.fileno()), "wb", buffer_size)
kk=out_ser.load_stream(infile)
foriteminkk:
print(item)

Python Server side:

importosimportpandasaspdos.environ["ARROW_PRE_0_15_IPC_FORMAT"] ="1"frompyjava.api.serveimportOnceServerddata=pd.DataFrame(data=[[1, 2, 3, 4], [2, 3, 4, 5]])
server=OnceServer("127.0.0.1", 11111, "Asia/Harbin")
server.bind()
server.serve([{'id': 9, 'label': 1}])

Java Client side:

importorg.apache.spark.sql.Rowimportorg.apache.spark.sql.catalyst.encoders.RowEncoderimportorg.apache.spark.sql.types.{LongType, StringType, StructField, StructType}
importorg.scalatest.{BeforeAndAfterAll, FunSuite}
importtech.mlsql.arrow.python.iapp.{AppContextImpl, JavaContext}
importtech.mlsql.arrow.python.runner.SparkSocketRunnerimporttech.mlsql.common.utils.network.NetUtilsvalenconder=RowEncoder.apply(StructType(Seq(StructField("a", LongType),StructField("b", LongType)))).resolveAndBind()
valsocketRunner=newSparkSocketRunner("wow", NetUtils.getHost, "Asia/Harbin")
valjavaConext=newJavaContextvalcommonTaskContext=newAppContextImpl(javaConext, null)
valiter= socketRunner.readFromStreamWithArrow("127.0.0.1", 11111, commonTaskContext)
iter.foreach(i => println(enconder.fromRow(i.copy())))
javaConext.close

How to configure python worker runs in Docker (todo)

About

No description, website, or topics provided.

Resources

Code of conduct

Contributing

Security policy

Stars

9 stars

Watchers

9 watching

Forks

Releases

Packages

Used by

Contributors

Languages